clientTmq.c 83.5 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
#include "tdef.h"
#include "tglobal.h"
X
Xiaoyu Wang 已提交
23
#include "tqueue.h"
24
#include "tref.h"
L
Liu Jicong 已提交
25 26
#include "ttimer.h"

X
Xiaoyu Wang 已提交
27 28
#define EMPTY_BLOCK_POLL_IDLE_DURATION 10
#define DEFAULT_AUTO_COMMIT_INTERVAL   5000
29

30 31
typedef void (*__tmq_askep_fn_t)(tmq_t* pTmq, int32_t code, SDataBuf* pBuf, void* pParam);

X
Xiaoyu Wang 已提交
32
struct SMqMgmt {
33 34 35
  int8_t  inited;
  tmr_h   timer;
  int32_t rsetId;
36
};
L
Liu Jicong 已提交
37

X
Xiaoyu Wang 已提交
38 39
static TdThreadOnce   tmqInit = PTHREAD_ONCE_INIT;  // initialize only once
volatile int32_t      tmqInitRes = 0;               // initialize rsp code
40
static struct SMqMgmt tmqMgmt = {0};
41

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

typedef struct {
L
Liu Jicong 已提交
48 49 50
  int8_t      tmqRspType;
  int32_t     epoch;
  SMqAskEpRsp msg;
L
Liu Jicong 已提交
51 52
} SMqAskEpRspWrapper;

L
Liu Jicong 已提交
53
struct tmq_list_t {
L
Liu Jicong 已提交
54
  SArray container;
L
Liu Jicong 已提交
55
};
L
Liu Jicong 已提交
56

L
Liu Jicong 已提交
57
struct tmq_conf_t {
58 59 60 61 62 63 64 65
  char           clientId[256];
  char           groupId[TSDB_CGROUP_LEN];
  int8_t         autoCommit;
  int8_t         resetOffset;
  int8_t         withTbName;
  int8_t         snapEnable;
  int32_t        snapBatchSize;
  bool           hbBgEnable;
66 67 68 69 70
  uint16_t       port;
  int32_t        autoCommitInterval;
  char*          ip;
  char*          user;
  char*          pass;
71
  tmq_commit_cb* commitCb;
L
Liu Jicong 已提交
72
  void*          commitCbUserParam;
L
Liu Jicong 已提交
73 74 75
};

struct tmq_t {
76 77 78 79 80 81 82 83 84 85
  int64_t        refId;
  char           groupId[TSDB_CGROUP_LEN];
  char           clientId[256];
  int8_t         withTbName;
  int8_t         useSnapshot;
  int8_t         autoCommit;
  int32_t        autoCommitInterval;
  int32_t        resetOffsetCfg;
  uint64_t       consumerId;
  bool           hbBgEnable;
L
Liu Jicong 已提交
86 87
  tmq_commit_cb* commitCb;
  void*          commitCbUserParam;
L
Liu Jicong 已提交
88 89 90 91

  // status
  int8_t  status;
  int32_t epoch;
L
Liu Jicong 已提交
92 93
#if 0
  int8_t  epStatus;
L
Liu Jicong 已提交
94
  int32_t epSkipCnt;
L
Liu Jicong 已提交
95
#endif
96
  // poll info
X
Xiaoyu Wang 已提交
97 98
  int64_t pollCnt;
  int64_t totalRows;
L
Liu Jicong 已提交
99

L
Liu Jicong 已提交
100
  // timer
X
Xiaoyu Wang 已提交
101 102 103 104 105 106 107 108 109 110
  tmr_h       hbLiveTimer;
  tmr_h       epTimer;
  tmr_h       reportTimer;
  tmr_h       commitTimer;
  STscObj*    pTscObj;       // connection
  SArray*     clientTopics;  // SArray<SMqClientTopic>
  STaosQueue* mqueue;        // queue of rsp
  STaosQall*  qall;
  STaosQueue* delayedTask;  // delayed task queue for heartbeat and auto commit
  tsem_t      rspSem;
L
Liu Jicong 已提交
111 112
};

113 114
typedef struct SAskEpInfo {
  int32_t code;
H
Haojun Liao 已提交
115
  tsem_t  sem;
116 117
} SAskEpInfo;

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

enum {
  TMQ_CONSUMER_STATUS__INIT = 0,
  TMQ_CONSUMER_STATUS__READY,
126
  TMQ_CONSUMER_STATUS__NO_TOPIC,
L
Liu Jicong 已提交
127
  TMQ_CONSUMER_STATUS__RECOVER,
L
Liu Jicong 已提交
128 129
};

L
Liu Jicong 已提交
130
enum {
131
  TMQ_DELAYED_TASK__ASK_EP = 1,
L
Liu Jicong 已提交
132 133 134 135
  TMQ_DELAYED_TASK__REPORT,
  TMQ_DELAYED_TASK__COMMIT,
};

H
Haojun Liao 已提交
136
typedef struct SVgOffsetInfo {
L
Liu Jicong 已提交
137 138
  STqOffsetVal committedOffset;
  STqOffsetVal currentOffset;
H
Haojun Liao 已提交
139 140 141 142 143 144 145 146 147 148 149
  int64_t      walVerBegin;
  int64_t      walVerEnd;
} SVgOffsetInfo;

typedef struct {
  int64_t       pollCnt;
  int64_t       numOfRows;
  SVgOffsetInfo offsetInfo;
  int32_t       vgId;
  int32_t       vgStatus;
  int32_t       vgSkipCnt;            // here used to mark the slow vgroups
150
  bool          receiveInfo;
H
Haojun Liao 已提交
151 152
  int64_t       emptyBlockReceiveTs;  // once empty block is received, idle for ignoreCnt then start to poll data
  SEpSet        epSet;
153 154
} SMqClientVg;

L
Liu Jicong 已提交
155
typedef struct {
156 157 158
  char           topicName[TSDB_TOPIC_FNAME_LEN];
  char           db[TSDB_DB_FNAME_LEN];
  SArray*        vgs;  // SArray<SMqClientVg>
L
Liu Jicong 已提交
159
  SSchemaWrapper schema;
160 161
} SMqClientTopic;

L
Liu Jicong 已提交
162 163
typedef struct {
  int8_t          tmqRspType;
164
  int32_t         epoch;  // epoch can be used to guard the vgHandle
165
  int32_t         vgId;
L
Liu Jicong 已提交
166 167
  SMqClientVg*    vgHandle;
  SMqClientTopic* topicHandle;
H
Haojun Liao 已提交
168
  uint64_t        reqId;
169
  SEpSet*         pEpset;
L
Liu Jicong 已提交
170
  union {
L
Liu Jicong 已提交
171 172
    SMqDataRsp dataRsp;
    SMqMetaRsp metaRsp;
L
Liu Jicong 已提交
173
    STaosxRsp  taosxRsp;
L
Liu Jicong 已提交
174
  };
L
Liu Jicong 已提交
175 176
} SMqPollRspWrapper;

L
Liu Jicong 已提交
177
typedef struct {
178 179
  int64_t refId;
  int32_t epoch;
L
Liu Jicong 已提交
180 181
  tsem_t  rspSem;
  int32_t rspErr;
L
Liu Jicong 已提交
182
} SMqSubscribeCbParam;
L
Liu Jicong 已提交
183

L
Liu Jicong 已提交
184
typedef struct {
185 186 187 188
  int64_t          refId;
  int32_t          epoch;
  void*            pParam;
  __tmq_askep_fn_t pUserFn;
189 190
} SMqAskEpCbParam;

L
Liu Jicong 已提交
191
typedef struct {
192 193
  int64_t         refId;
  int32_t         epoch;
L
Liu Jicong 已提交
194
  SMqClientVg*    pVg;
L
Liu Jicong 已提交
195
  SMqClientTopic* pTopic;
L
Liu Jicong 已提交
196
  int32_t         vgId;
X
Xiaoyu Wang 已提交
197
  uint64_t        requestId;  // request id for debug purpose
X
Xiaoyu Wang 已提交
198
} SMqPollCbParam;
199

200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216
typedef struct SMqVgCommon {
  tsem_t        rsp;
  int32_t       numOfRsp;
  SArray*       pList;
  TdThreadMutex mutex;
  int64_t       consumerId;
  char*         pTopicName;
  int32_t       code;
} SMqVgCommon;

typedef struct SMqVgWalInfoParam {
  int32_t      vgId;
  int32_t      epoch;
  int32_t      totalReq;
  SMqVgCommon* pCommon;
} SMqVgWalInfoParam;

217
typedef struct {
218 219
  int64_t        refId;
  int32_t        epoch;
L
Liu Jicong 已提交
220 221
  int32_t        waitingRspNum;
  int32_t        totalRspNum;
222
  int32_t        code;
223
  tmq_commit_cb* callbackFn;
L
Liu Jicong 已提交
224 225
  /*SArray*        successfulOffsets;*/
  /*SArray*        failedOffsets;*/
X
Xiaoyu Wang 已提交
226
  void* userParam;
227 228 229 230 231
} SMqCommitCbParamSet;

typedef struct {
  SMqCommitCbParamSet* params;
  STqOffset*           pOffset;
H
Haojun Liao 已提交
232 233 234
  char                 topicName[TSDB_TOPIC_FNAME_LEN];
  int32_t              vgId;
  tmq_t*               pTmq;
235
} SMqCommitCbParam;
236

237 238 239 240 241
typedef struct SSyncCommitInfo {
  tsem_t  sem;
  int32_t code;
} SSyncCommitInfo;

242
static int32_t doAskEp(tmq_t* tmq);
243 244
static int32_t makeTopicVgroupKey(char* dst, const char* topicName, int32_t vg);
static int32_t tmqCommitDone(SMqCommitCbParamSet* pParamSet);
245
static int32_t doSendCommitMsg(tmq_t* tmq, SMqClientVg* pVg, const char* pTopicName, SMqCommitCbParamSet* pParamSet,
246
                               int32_t index, int32_t totalVgroups, int32_t type);
X
Xiaoyu Wang 已提交
247 248 249
static void    commitRspCountDown(SMqCommitCbParamSet* pParamSet, int64_t consumerId, const char* pTopic, int32_t vgId);
static void    asyncAskEp(tmq_t* pTmq, __tmq_askep_fn_t askEpFn, void* param);
static void    addToQueueCallbackFn(tmq_t* pTmq, int32_t code, SDataBuf* pDataBuf, void* param);
250

251
tmq_conf_t* tmq_conf_new() {
wafwerar's avatar
wafwerar 已提交
252
  tmq_conf_t* conf = taosMemoryCalloc(1, sizeof(tmq_conf_t));
253 254 255 256 257
  if (conf == NULL) {
    terrno = TSDB_CODE_OUT_OF_MEMORY;
    return conf;
  }

258
  conf->withTbName = false;
L
Liu Jicong 已提交
259
  conf->autoCommit = true;
260
  conf->autoCommitInterval = DEFAULT_AUTO_COMMIT_INTERVAL;
261
  conf->resetOffset = TMQ_OFFSET__RESET_EARLIEAST;
262
  conf->hbBgEnable = true;
263

264 265 266
  return conf;
}

L
Liu Jicong 已提交
267
void tmq_conf_destroy(tmq_conf_t* conf) {
L
Liu Jicong 已提交
268
  if (conf) {
269 270 271 272 273 274 275 276 277
    if (conf->ip) {
      taosMemoryFree(conf->ip);
    }
    if (conf->user) {
      taosMemoryFree(conf->user);
    }
    if (conf->pass) {
      taosMemoryFree(conf->pass);
    }
L
Liu Jicong 已提交
278 279
    taosMemoryFree(conf);
  }
L
Liu Jicong 已提交
280 281 282
}

tmq_conf_res_t tmq_conf_set(tmq_conf_t* conf, const char* key, const char* value) {
283
  if (strcasecmp(key, "group.id") == 0) {
L
Liu Jicong 已提交
284
    tstrncpy(conf->groupId, value, TSDB_CGROUP_LEN);
L
Liu Jicong 已提交
285
    return TMQ_CONF_OK;
286
  }
L
Liu Jicong 已提交
287

288
  if (strcasecmp(key, "client.id") == 0) {
L
Liu Jicong 已提交
289
    tstrncpy(conf->clientId, value, 256);
L
Liu Jicong 已提交
290 291
    return TMQ_CONF_OK;
  }
L
Liu Jicong 已提交
292

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

305
  if (strcasecmp(key, "auto.commit.interval.ms") == 0) {
306
    conf->autoCommitInterval = taosStr2int64(value);
L
Liu Jicong 已提交
307 308 309
    return TMQ_CONF_OK;
  }

310 311 312
  if (strcasecmp(key, "auto.offset.reset") == 0) {
    if (strcasecmp(value, "none") == 0) {
      conf->resetOffset = TMQ_OFFSET__RESET_NONE;
L
Liu Jicong 已提交
313
      return TMQ_CONF_OK;
314 315
    } else if (strcasecmp(value, "earliest") == 0) {
      conf->resetOffset = TMQ_OFFSET__RESET_EARLIEAST;
L
Liu Jicong 已提交
316
      return TMQ_CONF_OK;
317 318
    } else if (strcasecmp(value, "latest") == 0) {
      conf->resetOffset = TMQ_OFFSET__RESET_LATEST;
L
Liu Jicong 已提交
319 320 321 322 323
      return TMQ_CONF_OK;
    } else {
      return TMQ_CONF_INVALID;
    }
  }
L
Liu Jicong 已提交
324

325 326
  if (strcasecmp(key, "msg.with.table.name") == 0) {
    if (strcasecmp(value, "true") == 0) {
327
      conf->withTbName = true;
L
Liu Jicong 已提交
328
      return TMQ_CONF_OK;
329
    } else if (strcasecmp(value, "false") == 0) {
330
      conf->withTbName = false;
L
Liu Jicong 已提交
331
      return TMQ_CONF_OK;
332 333 334 335 336
    } else {
      return TMQ_CONF_INVALID;
    }
  }

337 338
  if (strcasecmp(key, "experimental.snapshot.enable") == 0) {
    if (strcasecmp(value, "true") == 0) {
L
Liu Jicong 已提交
339
      conf->snapEnable = true;
340
      return TMQ_CONF_OK;
341
    } else if (strcasecmp(value, "false") == 0) {
L
Liu Jicong 已提交
342
      conf->snapEnable = false;
343 344 345 346 347 348
      return TMQ_CONF_OK;
    } else {
      return TMQ_CONF_INVALID;
    }
  }

349
  if (strcasecmp(key, "experimental.snapshot.batch.size") == 0) {
350
    conf->snapBatchSize = taosStr2int64(value);
L
Liu Jicong 已提交
351 352 353
    return TMQ_CONF_OK;
  }

354
  if (strcasecmp(key, "enable.heartbeat.background") == 0) {
X
Xiaoyu Wang 已提交
355 356 357 358 359 360 361 362 363 364
    //    if (strcasecmp(value, "true") == 0) {
    //      conf->hbBgEnable = true;
    //      return TMQ_CONF_OK;
    //    } else if (strcasecmp(value, "false") == 0) {
    //      conf->hbBgEnable = false;
    //      return TMQ_CONF_OK;
    //    } else {
    tscError("the default value of enable.heartbeat.background is true, can not be seted");
    return TMQ_CONF_INVALID;
    //    }
L
Liu Jicong 已提交
365 366
  }

367
  if (strcasecmp(key, "td.connect.ip") == 0) {
368
    conf->ip = taosStrdup(value);
L
Liu Jicong 已提交
369 370
    return TMQ_CONF_OK;
  }
371

372
  if (strcasecmp(key, "td.connect.user") == 0) {
373
    conf->user = taosStrdup(value);
L
Liu Jicong 已提交
374 375
    return TMQ_CONF_OK;
  }
376

377
  if (strcasecmp(key, "td.connect.pass") == 0) {
378
    conf->pass = taosStrdup(value);
L
Liu Jicong 已提交
379 380
    return TMQ_CONF_OK;
  }
381

382
  if (strcasecmp(key, "td.connect.port") == 0) {
383
    conf->port = taosStr2int64(value);
L
Liu Jicong 已提交
384 385
    return TMQ_CONF_OK;
  }
386

387
  if (strcasecmp(key, "td.connect.db") == 0) {
L
Liu Jicong 已提交
388 389 390
    return TMQ_CONF_OK;
  }

L
Liu Jicong 已提交
391
  return TMQ_CONF_UNKNOWN;
392 393
}

X
Xiaoyu Wang 已提交
394
tmq_list_t* tmq_list_new() { return (tmq_list_t*)taosArrayInit(0, sizeof(void*)); }
395

L
Liu Jicong 已提交
396 397
int32_t tmq_list_append(tmq_list_t* list, const char* src) {
  SArray* container = &list->container;
398
  if (src == NULL || src[0] == 0) return -1;
399
  char* topic = taosStrdup(src);
400 401 402
  if (topic[0] != '`') {
    strtolower(topic, src);
  }
L
fix  
Liu Jicong 已提交
403
  if (taosArrayPush(container, &topic) == NULL) return -1;
404 405 406
  return 0;
}

L
Liu Jicong 已提交
407
void tmq_list_destroy(tmq_list_t* list) {
L
Liu Jicong 已提交
408
  SArray* container = &list->container;
L
Liu Jicong 已提交
409
  taosArrayDestroyP(container, taosMemoryFree);
L
Liu Jicong 已提交
410 411
}

L
Liu Jicong 已提交
412 413 414 415 416 417 418 419 420 421
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;
}

X
Xiaoyu Wang 已提交
422 423
static SMqClientVg* foundClientVg(SArray* pTopicList, const char* pName, int32_t vgId, int32_t* index,
                                  int32_t* numOfVgroups) {
424 425 426 427
  int32_t numOfTopics = taosArrayGetSize(pTopicList);
  *index = -1;
  *numOfVgroups = 0;

X
Xiaoyu Wang 已提交
428
  for (int32_t i = 0; i < numOfTopics; ++i) {
429 430
    SMqClientTopic* pTopic = taosArrayGet(pTopicList, i);
    if (strcmp(pTopic->topicName, pName) != 0) {
431 432 433
      continue;
    }

434 435
    *numOfVgroups = taosArrayGetSize(pTopic->vgs);
    for (int32_t j = 0; j < (*numOfVgroups); ++j) {
436
      SMqClientVg* pClientVg = taosArrayGet(pTopic->vgs, j);
437 438 439
      if (pClientVg->vgId == vgId) {
        *index = j;
        return pClientVg;
440 441
      }
    }
L
Liu Jicong 已提交
442
  }
443 444

  return NULL;
L
Liu Jicong 已提交
445
}
446

447 448 449
// Two problems do not need to be addressed here
// 1. update to of epset. the response of poll request will automatically handle this problem
// 2. commit failure. This one needs to be resolved.
H
Haojun Liao 已提交
450
static int32_t tmqCommitCb(void* param, SDataBuf* pBuf, int32_t code) {
451
  SMqCommitCbParam*    pParam = (SMqCommitCbParam*)param;
452
  SMqCommitCbParamSet* pParamSet = (SMqCommitCbParamSet*)pParam->params;
H
Haojun Liao 已提交
453

X
Xiaoyu Wang 已提交
454 455 456 457 458 459 460 461 462 463 464 465 466 467
  //  if (code != TSDB_CODE_SUCCESS) { // if commit offset failed, let's try again
  //    taosThreadMutexLock(&pParam->pTmq->lock);
  //    int32_t numOfVgroups, index;
  //    SMqClientVg* pVg = foundClientVg(pParam->pTmq->clientTopics, pParam->topicName, pParam->vgId, &index,
  //    &numOfVgroups); if (pVg == NULL) {
  //      tscDebug("consumer:0x%" PRIx64
  //               " subKey:%s vgId:%d commit failed, code:%s has been transferred to other consumer, no need retry
  //               ordinal:%d/%d", pParam->pTmq->consumerId, pParam->pOffset->subKey, pParam->vgId, tstrerror(code),
  //               index + 1, numOfVgroups);
  //    } else { // let's retry the commit
  //      int32_t code1 = doSendCommitMsg(pParam->pTmq, pVg, pParam->topicName, pParamSet, index, numOfVgroups);
  //      if (code1 != TSDB_CODE_SUCCESS) {  // retry failed.
  //        tscError("consumer:0x%" PRIx64 " topic:%s vgId:%d offset:%" PRId64
  //                 " retry failed, ignore this commit. code:%s ordinal:%d/%d",
H
Haojun Liao 已提交
468
  //                 pParam->pTmq->consumerId, pParam->topicName, pVg->vgId, pVg->offsetInfo.committedOffset.version,
X
Xiaoyu Wang 已提交
469 470 471 472 473 474 475 476 477 478 479 480 481 482 483
  //                 tstrerror(terrno), index + 1, numOfVgroups);
  //      }
  //    }
  //
  //    taosThreadMutexUnlock(&pParam->pTmq->lock);
  //
  //    taosMemoryFree(pParam->pOffset);
  //    taosMemoryFree(pBuf->pData);
  //    taosMemoryFree(pBuf->pEpSet);
  //
  //    commitRspCountDown(pParamSet, pParam->pTmq->consumerId, pParam->topicName, pParam->vgId);
  //    return 0;
  //  }
  //
  //  // todo replace the pTmq with refId
484

L
Liu Jicong 已提交
485
  taosMemoryFree(pParam->pOffset);
L
Liu Jicong 已提交
486
  taosMemoryFree(pBuf->pData);
dengyihao's avatar
dengyihao 已提交
487
  taosMemoryFree(pBuf->pEpSet);
L
Liu Jicong 已提交
488

489
  commitRspCountDown(pParamSet, pParam->pTmq->consumerId, pParam->topicName, pParam->vgId);
490 491 492
  return 0;
}

493
static int32_t doSendCommitMsg(tmq_t* tmq, SMqClientVg* pVg, const char* pTopicName, SMqCommitCbParamSet* pParamSet,
494
                               int32_t index, int32_t totalVgroups, int32_t type) {
L
Liu Jicong 已提交
495 496
  STqOffset* pOffset = taosMemoryCalloc(1, sizeof(STqOffset));
  if (pOffset == NULL) {
497
    return TSDB_CODE_OUT_OF_MEMORY;
L
Liu Jicong 已提交
498
  }
499

H
Haojun Liao 已提交
500
  pOffset->val = pVg->offsetInfo.currentOffset;
501

L
Liu Jicong 已提交
502 503 504
  int32_t groupLen = strlen(tmq->groupId);
  memcpy(pOffset->subKey, tmq->groupId, groupLen);
  pOffset->subKey[groupLen] = TMQ_SEPARATOR;
H
Haojun Liao 已提交
505
  strcpy(pOffset->subKey + groupLen + 1, pTopicName);
L
Liu Jicong 已提交
506

507 508
  int32_t len = 0;
  int32_t code = 0;
L
Liu Jicong 已提交
509 510
  tEncodeSize(tEncodeSTqOffset, pOffset, len, code);
  if (code < 0) {
511
    return TSDB_CODE_INVALID_PARA;
L
Liu Jicong 已提交
512
  }
513

L
Liu Jicong 已提交
514
  void* buf = taosMemoryCalloc(1, sizeof(SMsgHead) + len);
L
Liu Jicong 已提交
515 516
  if (buf == NULL) {
    taosMemoryFree(pOffset);
517
    return TSDB_CODE_OUT_OF_MEMORY;
L
Liu Jicong 已提交
518
  }
519

L
Liu Jicong 已提交
520
  ((SMsgHead*)buf)->vgId = htonl(pVg->vgId);
L
Liu Jicong 已提交
521

L
Liu Jicong 已提交
522 523 524 525 526
  void* abuf = POINTER_SHIFT(buf, sizeof(SMsgHead));

  SEncoder encoder;
  tEncoderInit(&encoder, abuf, len);
  tEncodeSTqOffset(&encoder, pOffset);
L
Liu Jicong 已提交
527
  tEncoderClear(&encoder);
L
Liu Jicong 已提交
528 529

  // build param
530
  SMqCommitCbParam* pParam = taosMemoryCalloc(1, sizeof(SMqCommitCbParam));
L
Liu Jicong 已提交
531
  if (pParam == NULL) {
L
Liu Jicong 已提交
532
    taosMemoryFree(pOffset);
L
Liu Jicong 已提交
533
    taosMemoryFree(buf);
534
    return TSDB_CODE_OUT_OF_MEMORY;
L
Liu Jicong 已提交
535
  }
536

L
Liu Jicong 已提交
537 538
  pParam->params = pParamSet;
  pParam->pOffset = pOffset;
H
Haojun Liao 已提交
539 540 541
  pParam->vgId = pVg->vgId;
  pParam->pTmq = tmq;

H
Haojun Liao 已提交
542
  tstrncpy(pParam->topicName, pTopicName, tListLen(pParam->topicName));
L
Liu Jicong 已提交
543 544 545 546

  // build send info
  SMsgSendInfo* pMsgSendInfo = taosMemoryCalloc(1, sizeof(SMsgSendInfo));
  if (pMsgSendInfo == NULL) {
L
Liu Jicong 已提交
547
    taosMemoryFree(pOffset);
L
Liu Jicong 已提交
548 549
    taosMemoryFree(buf);
    taosMemoryFree(pParam);
550
    return TSDB_CODE_OUT_OF_MEMORY;
L
Liu Jicong 已提交
551
  }
552

553
  pMsgSendInfo->msgInfo = (SDataBuf) { .pData = buf, .len = sizeof(SMsgHead) + len, .handle = NULL };
L
Liu Jicong 已提交
554 555 556 557

  pMsgSendInfo->requestId = generateRequestId();
  pMsgSendInfo->requestObjRefId = 0;
  pMsgSendInfo->param = pParam;
L
Liu Jicong 已提交
558
  pMsgSendInfo->paramFreeFp = taosMemoryFree;
559
  pMsgSendInfo->fp = tmqCommitCb;
560
  pMsgSendInfo->msgType = type;
L
Liu Jicong 已提交
561

L
Liu Jicong 已提交
562 563 564
  atomic_add_fetch_32(&pParamSet->waitingRspNum, 1);
  atomic_add_fetch_32(&pParamSet->totalRspNum, 1);

H
Haojun Liao 已提交
565
  SEp* pEp = GET_ACTIVE_EP(&pVg->epSet);
566 567 568 569
  char offsetBuf[80] = {0};
  tFormatOffset(offsetBuf, tListLen(offsetBuf), &pOffset->val);

  char commitBuf[80] = {0};
H
Haojun Liao 已提交
570
  tFormatOffset(commitBuf, tListLen(commitBuf), &pVg->offsetInfo.committedOffset);
571 572 573
  tscDebug("consumer:0x%" PRIx64 " topic:%s on vgId:%d send offset:%s prev:%s, ep:%s:%d, ordinal:%d/%d, req:0x%" PRIx64,
           tmq->consumerId, pOffset->subKey, pVg->vgId, offsetBuf, commitBuf, pEp->fqdn, pEp->port, index + 1,
           totalVgroups, pMsgSendInfo->requestId);
H
Haojun Liao 已提交
574

L
Liu Jicong 已提交
575 576
  int64_t transporterId = 0;
  asyncSendMsgToServer(tmq->pTscObj->pAppInfo->pTransporter, &pVg->epSet, &transporterId, pMsgSendInfo);
577 578

  return TSDB_CODE_SUCCESS;
L
Liu Jicong 已提交
579 580
}

H
Haojun Liao 已提交
581 582 583 584 585 586 587 588 589 590 591
static SMqClientTopic* getTopicByName(tmq_t* tmq, const char* pTopicName) {
  int32_t numOfTopics = taosArrayGetSize(tmq->clientTopics);
  for (int32_t i = 0; i < numOfTopics; ++i) {
    SMqClientTopic* pTopic = taosArrayGet(tmq->clientTopics, i);
    if (strcmp(pTopic->topicName, pTopicName) != 0) {
      continue;
    }

    return pTopic;
  }

H
Haojun Liao 已提交
592
  tscError("consumer:0x%" PRIx64 ", total:%d, failed to find topic:%s", tmq->consumerId, numOfTopics, pTopicName);
H
Haojun Liao 已提交
593 594 595
  return NULL;
}

596
static void asyncCommitOffset(tmq_t* tmq, const TAOS_RES* pRes, int32_t type, tmq_commit_cb* pCommitFp, void* userParam) {
597 598 599 600 601 602 603 604 605 606 607 608
  char*   pTopicName = NULL;
  int32_t vgId = 0;
  int32_t code = 0;

  if (pRes == NULL || tmq == NULL) {
    pCommitFp(tmq, TSDB_CODE_INVALID_PARA, userParam);
    return;
  }

  if (TD_RES_TMQ(pRes)) {
    SMqRspObj* pRspObj = (SMqRspObj*)pRes;
    pTopicName = pRspObj->topic;
L
Liu Jicong 已提交
609
    vgId = pRspObj->vgId;
610 611 612
  } else if (TD_RES_TMQ_META(pRes)) {
    SMqMetaRspObj* pMetaRspObj = (SMqMetaRspObj*)pRes;
    pTopicName = pMetaRspObj->topic;
L
Liu Jicong 已提交
613
    vgId = pMetaRspObj->vgId;
614 615 616
  } else if (TD_RES_TMQ_METADATA(pRes)) {
    SMqTaosxRspObj* pRspObj = (SMqTaosxRspObj*)pRes;
    pTopicName = pRspObj->topic;
617
    vgId = pRspObj->vgId;
L
Liu Jicong 已提交
618
  } else {
619 620
    pCommitFp(tmq, TSDB_CODE_TMQ_INVALID_MSG, userParam);
    return;
L
Liu Jicong 已提交
621 622 623 624
  }

  SMqCommitCbParamSet* pParamSet = taosMemoryCalloc(1, sizeof(SMqCommitCbParamSet));
  if (pParamSet == NULL) {
625 626
    pCommitFp(tmq, TSDB_CODE_OUT_OF_MEMORY, userParam);
    return;
L
Liu Jicong 已提交
627
  }
H
Haojun Liao 已提交
628

629 630
  pParamSet->refId = tmq->refId;
  pParamSet->epoch = tmq->epoch;
631
  pParamSet->callbackFn = pCommitFp;
L
Liu Jicong 已提交
632
  pParamSet->userParam = userParam;
L
Liu Jicong 已提交
633

H
Haojun Liao 已提交
634 635
  int32_t numOfTopics = taosArrayGetSize(tmq->clientTopics);

636 637
  tscDebug("consumer:0x%" PRIx64 " do manual commit offset for %s, vgId:%d", tmq->consumerId, pTopicName, vgId);

H
Haojun Liao 已提交
638 639
  SMqClientTopic* pTopic = getTopicByName(tmq, pTopicName);
  if (pTopic == NULL) {
X
Xiaoyu Wang 已提交
640 641
    tscWarn("consumer:0x%" PRIx64 " failed to find the specified topic:%s, total topics:%d", tmq->consumerId,
            pTopicName, numOfTopics);
642 643 644 645
    taosMemoryFree(pParamSet);
    pCommitFp(tmq, TSDB_CODE_SUCCESS, userParam);
    return;
  }
L
Liu Jicong 已提交
646

647 648 649 650 651 652
  int32_t j = 0;
  int32_t numOfVgroups = taosArrayGetSize(pTopic->vgs);
  for (j = 0; j < numOfVgroups; j++) {
    SMqClientVg* pVg = taosArrayGet(pTopic->vgs, j);
    if (pVg->vgId == vgId) {
      break;
L
Liu Jicong 已提交
653
    }
L
Liu Jicong 已提交
654
  }
L
Liu Jicong 已提交
655

656
  if (j == numOfVgroups) {
X
Xiaoyu Wang 已提交
657 658
    tscWarn("consumer:0x%" PRIx64 " failed to find the specified vgId:%d, total Vgs:%d, topic:%s", tmq->consumerId,
            vgId, numOfVgroups, pTopicName);
L
Liu Jicong 已提交
659
    taosMemoryFree(pParamSet);
660 661
    pCommitFp(tmq, TSDB_CODE_SUCCESS, userParam);
    return;
L
Liu Jicong 已提交
662 663
  }

664
  SMqClientVg* pVg = taosArrayGet(pTopic->vgs, j);
H
Haojun Liao 已提交
665
  if (pVg->offsetInfo.currentOffset.type > 0 && !tOffsetEqual(&pVg->offsetInfo.currentOffset, &pVg->offsetInfo.committedOffset)) {
666
    code = doSendCommitMsg(tmq, pVg, pTopic->topicName, pParamSet, j, numOfVgroups, type);
L
Liu Jicong 已提交
667

668 669 670 671 672
    // failed to commit, callback user function directly.
    if (code != TSDB_CODE_SUCCESS) {
      taosMemoryFree(pParamSet);
      pCommitFp(tmq, code, userParam);
    }
X
Xiaoyu Wang 已提交
673
  } else {  // do not perform commit, callback user function directly.
674 675
    taosMemoryFree(pParamSet);
    pCommitFp(tmq, code, userParam);
L
Liu Jicong 已提交
676 677 678
  }
}

679
static void asyncCommitAllOffsets(tmq_t* tmq, tmq_commit_cb* pCommitFp, void* userParam) {
680 681
  SMqCommitCbParamSet* pParamSet = taosMemoryCalloc(1, sizeof(SMqCommitCbParamSet));
  if (pParamSet == NULL) {
682 683
    pCommitFp(tmq, TSDB_CODE_OUT_OF_MEMORY, userParam);
    return;
684
  }
685 686 687

  pParamSet->refId = tmq->refId;
  pParamSet->epoch = tmq->epoch;
688
  pParamSet->callbackFn = pCommitFp;
689 690
  pParamSet->userParam = userParam;

691 692 693
  // init as 1 to prevent concurrency issue
  pParamSet->waitingRspNum = 1;

694
  int32_t numOfTopics = taosArrayGetSize(tmq->clientTopics);
X
Xiaoyu Wang 已提交
695
  tscDebug("consumer:0x%" PRIx64 " start to commit offset for %d topics", tmq->consumerId, numOfTopics);
696 697

  for (int32_t i = 0; i < numOfTopics; i++) {
698
    SMqClientTopic* pTopic = taosArrayGet(tmq->clientTopics, i);
699
    int32_t         numOfVgroups = taosArrayGetSize(pTopic->vgs);
L
Liu Jicong 已提交
700

701 702
    tscDebug("consumer:0x%" PRIx64 " commit offset for topics:%s, numOfVgs:%d", tmq->consumerId, pTopic->topicName,
             numOfVgroups);
703
    for (int32_t j = 0; j < numOfVgroups; j++) {
704 705
      SMqClientVg* pVg = taosArrayGet(pTopic->vgs, j);

H
Haojun Liao 已提交
706
      if (pVg->offsetInfo.currentOffset.type > 0 && !tOffsetEqual(&pVg->offsetInfo.currentOffset, &pVg->offsetInfo.committedOffset)) {
707
        int32_t code = doSendCommitMsg(tmq, pVg, pTopic->topicName, pParamSet, j, numOfVgroups, TDMT_VND_TMQ_COMMIT_OFFSET);
708 709
        if (code != TSDB_CODE_SUCCESS) {
          tscError("consumer:0x%" PRIx64 " topic:%s vgId:%d offset:%" PRId64 " failed, code:%s ordinal:%d/%d",
H
Haojun Liao 已提交
710
                   tmq->consumerId, pTopic->topicName, pVg->vgId, pVg->offsetInfo.committedOffset.version, tstrerror(terrno),
711
                   j + 1, numOfVgroups);
L
Liu Jicong 已提交
712 713
          continue;
        }
H
Haojun Liao 已提交
714 715

        // update the offset value.
H
Haojun Liao 已提交
716
        pVg->offsetInfo.committedOffset = pVg->offsetInfo.currentOffset;
717
      } else {
D
dapan1121 已提交
718
        tscDebug("consumer:0x%" PRIx64 " topic:%s vgId:%d, no commit, current:%" PRId64 ", ordinal:%d/%d",
H
Haojun Liao 已提交
719
                 tmq->consumerId, pTopic->topicName, pVg->vgId, pVg->offsetInfo.currentOffset.version, j + 1, numOfVgroups);
720 721 722 723
      }
    }
  }

H
Haojun Liao 已提交
724
  tscDebug("consumer:0x%" PRIx64 " total commit:%d for %d topics", tmq->consumerId, pParamSet->waitingRspNum - 1,
725
           numOfTopics);
H
Haojun Liao 已提交
726

L
Liu Jicong 已提交
727
  // no request is sent
L
Liu Jicong 已提交
728 729
  if (pParamSet->totalRspNum == 0) {
    taosMemoryFree(pParamSet);
730 731
    pCommitFp(tmq, TSDB_CODE_SUCCESS, userParam);
    return;
L
Liu Jicong 已提交
732 733
  }

L
Liu Jicong 已提交
734
  // count down since waiting rsp num init as 1
735
  commitRspCountDown(pParamSet, tmq->consumerId, "", 0);
736 737
}

738 739
static void generateTimedTask(int64_t refId, int32_t type) {
  tmq_t* tmq = taosAcquireRef(tmqMgmt.rsetId, refId);
740
  if (tmq != NULL) {
S
Shengliang Guan 已提交
741
    int8_t* pTaskType = taosAllocateQitem(sizeof(int8_t), DEF_QITEM, 0);
742
    *pTaskType = type;
743 744 745
    taosWriteQitem(tmq->delayedTask, pTaskType);
    tsem_post(&tmq->rspSem);
  }
746
  taosReleaseRef(tmqMgmt.rsetId, refId);
747 748 749 750 751
}

void tmqAssignAskEpTask(void* param, void* tmrId) {
  int64_t refId = *(int64_t*)param;
  generateTimedTask(refId, TMQ_DELAYED_TASK__ASK_EP);
752
  taosMemoryFree(param);
L
Liu Jicong 已提交
753 754 755
}

void tmqAssignDelayedCommitTask(void* param, void* tmrId) {
756
  int64_t refId = *(int64_t*)param;
757
  generateTimedTask(refId, TMQ_DELAYED_TASK__COMMIT);
758
  taosMemoryFree(param);
L
Liu Jicong 已提交
759 760 761
}

void tmqAssignDelayedReportTask(void* param, void* tmrId) {
762 763 764
  int64_t refId = *(int64_t*)param;
  tmq_t*  tmq = taosAcquireRef(tmqMgmt.rsetId, refId);
  if (tmq != NULL) {
S
Shengliang Guan 已提交
765
    int8_t* pTaskType = taosAllocateQitem(sizeof(int8_t), DEF_QITEM, 0);
766 767 768 769
    *pTaskType = TMQ_DELAYED_TASK__REPORT;
    taosWriteQitem(tmq->delayedTask, pTaskType);
    tsem_post(&tmq->rspSem);
  }
770 771

  taosReleaseRef(tmqMgmt.rsetId, refId);
772
  taosMemoryFree(param);
L
Liu Jicong 已提交
773 774
}

775
int32_t tmqHbCb(void* param, SDataBuf* pMsg, int32_t code) {
dengyihao's avatar
dengyihao 已提交
776 777 778 779
  if (pMsg) {
    taosMemoryFree(pMsg->pData);
    taosMemoryFree(pMsg->pEpSet);
  }
780 781 782 783
  return 0;
}

void tmqSendHbReq(void* param, void* tmrId) {
784
  int64_t refId = *(int64_t*)param;
785

X
Xiaoyu Wang 已提交
786
  tmq_t* tmq = taosAcquireRef(tmqMgmt.rsetId, refId);
787
  if (tmq == NULL) {
L
Liu Jicong 已提交
788
    taosMemoryFree(param);
789 790
    return;
  }
D
dapan1121 已提交
791 792 793 794 795

  SMqHbReq req = {0};
  req.consumerId = tmq->consumerId;
  req.epoch = tmq->epoch;

L
Liu Jicong 已提交
796
  int32_t tlen = tSerializeSMqHbReq(NULL, 0, &req);
D
dapan1121 已提交
797 798
  if (tlen < 0) {
    tscError("tSerializeSMqHbReq failed");
799
    goto OVER;
D
dapan1121 已提交
800
  }
801

L
Liu Jicong 已提交
802
  void* pReq = taosMemoryCalloc(1, tlen);
D
dapan1121 已提交
803 804
  if (tlen < 0) {
    tscError("failed to malloc MqHbReq msg, size:%d", tlen);
805
    goto OVER;
D
dapan1121 已提交
806
  }
807

D
dapan1121 已提交
808 809 810
  if (tSerializeSMqHbReq(pReq, tlen, &req) < 0) {
    tscError("tSerializeSMqHbReq %d failed", tlen);
    taosMemoryFree(pReq);
811
    goto OVER;
D
dapan1121 已提交
812
  }
813 814 815 816

  SMsgSendInfo* sendInfo = taosMemoryCalloc(1, sizeof(SMsgSendInfo));
  if (sendInfo == NULL) {
    taosMemoryFree(pReq);
L
Liu Jicong 已提交
817
    goto OVER;
818
  }
819

820
  sendInfo->msgInfo = (SDataBuf){ .pData = pReq, .len = tlen, .handle = NULL };
821 822 823 824 825

  sendInfo->requestId = generateRequestId();
  sendInfo->requestObjRefId = 0;
  sendInfo->param = NULL;
  sendInfo->fp = tmqHbCb;
L
Liu Jicong 已提交
826
  sendInfo->msgType = TDMT_MND_TMQ_HB;
827 828 829 830 831 832 833

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

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

OVER:
834
  taosTmrReset(tmqSendHbReq, 1000, param, tmqMgmt.timer, &tmq->hbLiveTimer);
835
  taosReleaseRef(tmqMgmt.rsetId, refId);
836 837
}

838 839
static void defaultCommitCbFn(tmq_t* pTmq, int32_t code, void* param) {
  if (code != 0) {
X
Xiaoyu Wang 已提交
840
    tscDebug("consumer:0x%" PRIx64 ", failed to commit offset, code:%s", pTmq->consumerId, tstrerror(code));
841 842 843
  }
}

844
int32_t tmqHandleAllDelayedTask(tmq_t* pTmq) {
L
Liu Jicong 已提交
845
  STaosQall* qall = taosAllocateQall();
846
  taosReadAllQitems(pTmq->delayedTask, qall);
L
Liu Jicong 已提交
847

848 849 850 851
  if (qall->numOfItems == 0) {
    taosFreeQall(qall);
    return TSDB_CODE_SUCCESS;
  }
852

X
Xiaoyu Wang 已提交
853
  tscDebug("consumer:0x%" PRIx64 " handle delayed %d tasks before poll data", pTmq->consumerId, qall->numOfItems);
854 855
  int8_t* pTaskType = NULL;
  taosGetQitem(qall, (void**)&pTaskType);
L
Liu Jicong 已提交
856

857
  while (pTaskType != NULL) {
858
    if (*pTaskType == TMQ_DELAYED_TASK__ASK_EP) {
859
      asyncAskEp(pTmq, addToQueueCallbackFn, NULL);
860 861

      int64_t* pRefId = taosMemoryMalloc(sizeof(int64_t));
862
      *pRefId = pTmq->refId;
863

X
Xiaoyu Wang 已提交
864
      tscDebug("consumer:0x%" PRIx64 " retrieve ep from mnode in 1s", pTmq->consumerId);
865
      taosTmrReset(tmqAssignAskEpTask, 1000, pRefId, tmqMgmt.timer, &pTmq->epTimer);
L
Liu Jicong 已提交
866
    } else if (*pTaskType == TMQ_DELAYED_TASK__COMMIT) {
X
Xiaoyu Wang 已提交
867
      tmq_commit_cb* pCallbackFn = pTmq->commitCb ? pTmq->commitCb : defaultCommitCbFn;
868 869

      asyncCommitAllOffsets(pTmq, pCallbackFn, pTmq->commitCbUserParam);
870
      int64_t* pRefId = taosMemoryMalloc(sizeof(int64_t));
871
      *pRefId = pTmq->refId;
872

873
      tscDebug("consumer:0x%" PRIx64 " next commit to vnode(s) in %.2fs", pTmq->consumerId,
X
Xiaoyu Wang 已提交
874
               pTmq->autoCommitInterval / 1000.0);
875
      taosTmrReset(tmqAssignDelayedCommitTask, pTmq->autoCommitInterval, pRefId, tmqMgmt.timer, &pTmq->commitTimer);
L
Liu Jicong 已提交
876 877
    } else if (*pTaskType == TMQ_DELAYED_TASK__REPORT) {
    }
878

L
Liu Jicong 已提交
879
    taosFreeQitem(pTaskType);
880
    taosGetQitem(qall, (void**)&pTaskType);
L
Liu Jicong 已提交
881
  }
882

L
Liu Jicong 已提交
883 884 885 886
  taosFreeQall(qall);
  return 0;
}

887
static void* tmqFreeRspWrapper(SMqRspWrapper* rspWrapper) {
L
Liu Jicong 已提交
888 889 890 891 892 893 894
  if (rspWrapper->tmqRspType == TMQ_MSG_TYPE__END_RSP) {
    // do nothing
  } else if (rspWrapper->tmqRspType == TMQ_MSG_TYPE__EP_RSP) {
    SMqAskEpRspWrapper* pEpRspWrapper = (SMqAskEpRspWrapper*)rspWrapper;
    tDeleteSMqAskEpRsp(&pEpRspWrapper->msg);
  } else if (rspWrapper->tmqRspType == TMQ_MSG_TYPE__POLL_RSP) {
    SMqPollRspWrapper* pRsp = (SMqPollRspWrapper*)rspWrapper;
895 896
    taosMemoryFreeClear(pRsp->pEpset);

L
Liu Jicong 已提交
897 898 899 900 901 902
    taosArrayDestroyP(pRsp->dataRsp.blockData, taosMemoryFree);
    taosArrayDestroy(pRsp->dataRsp.blockDataLen);
    taosArrayDestroyP(pRsp->dataRsp.blockTbName, taosMemoryFree);
    taosArrayDestroyP(pRsp->dataRsp.blockSchema, (FDelete)tDeleteSSchemaWrapper);
  } else if (rspWrapper->tmqRspType == TMQ_MSG_TYPE__POLL_META_RSP) {
    SMqPollRspWrapper* pRsp = (SMqPollRspWrapper*)rspWrapper;
903 904
    taosMemoryFreeClear(pRsp->pEpset);

L
Liu Jicong 已提交
905 906 907
    taosMemoryFree(pRsp->metaRsp.metaRsp);
  } else if (rspWrapper->tmqRspType == TMQ_MSG_TYPE__TAOSX_RSP) {
    SMqPollRspWrapper* pRsp = (SMqPollRspWrapper*)rspWrapper;
908 909
    taosMemoryFreeClear(pRsp->pEpset);

L
Liu Jicong 已提交
910 911 912 913 914 915 916 917
    taosArrayDestroyP(pRsp->taosxRsp.blockData, taosMemoryFree);
    taosArrayDestroy(pRsp->taosxRsp.blockDataLen);
    taosArrayDestroyP(pRsp->taosxRsp.blockTbName, taosMemoryFree);
    taosArrayDestroyP(pRsp->taosxRsp.blockSchema, (FDelete)tDeleteSSchemaWrapper);
    // taosx
    taosArrayDestroy(pRsp->taosxRsp.createTableLen);
    taosArrayDestroyP(pRsp->taosxRsp.createTableReq, taosMemoryFree);
  }
918 919

  return NULL;
L
Liu Jicong 已提交
920 921
}

L
Liu Jicong 已提交
922
void tmqClearUnhandleMsg(tmq_t* tmq) {
L
Liu Jicong 已提交
923
  SMqRspWrapper* rspWrapper = NULL;
L
Liu Jicong 已提交
924
  while (1) {
L
Liu Jicong 已提交
925 926 927 928 929
    taosGetQitem(tmq->qall, (void**)&rspWrapper);
    if (rspWrapper) {
      tmqFreeRspWrapper(rspWrapper);
      taosFreeQitem(rspWrapper);
    } else {
L
Liu Jicong 已提交
930
      break;
L
Liu Jicong 已提交
931
    }
L
Liu Jicong 已提交
932 933
  }

L
Liu Jicong 已提交
934
  rspWrapper = NULL;
L
Liu Jicong 已提交
935 936
  taosReadAllQitems(tmq->mqueue, tmq->qall);
  while (1) {
L
Liu Jicong 已提交
937 938 939 940 941
    taosGetQitem(tmq->qall, (void**)&rspWrapper);
    if (rspWrapper) {
      tmqFreeRspWrapper(rspWrapper);
      taosFreeQitem(rspWrapper);
    } else {
L
Liu Jicong 已提交
942
      break;
L
Liu Jicong 已提交
943
    }
L
Liu Jicong 已提交
944 945 946
  }
}

D
dapan1121 已提交
947
int32_t tmqSubscribeCb(void* param, SDataBuf* pMsg, int32_t code) {
L
Liu Jicong 已提交
948 949
  SMqSubscribeCbParam* pParam = (SMqSubscribeCbParam*)param;
  pParam->rspErr = code;
dengyihao's avatar
dengyihao 已提交
950 951

  taosMemoryFree(pMsg->pEpSet);
L
Liu Jicong 已提交
952 953 954
  tsem_post(&pParam->rspSem);
  return 0;
}
955

L
Liu Jicong 已提交
956
int32_t tmq_subscription(tmq_t* tmq, tmq_list_t** topics) {
X
Xiaoyu Wang 已提交
957 958 959 960
  if (*topics == NULL) {
    *topics = tmq_list_new();
  }
  for (int i = 0; i < taosArrayGetSize(tmq->clientTopics); i++) {
L
Liu Jicong 已提交
961
    SMqClientTopic* topic = taosArrayGet(tmq->clientTopics, i);
L
Liu Jicong 已提交
962
    tmq_list_append(*topics, strchr(topic->topicName, '.') + 1);
X
Xiaoyu Wang 已提交
963
  }
L
Liu Jicong 已提交
964
  return 0;
X
Xiaoyu Wang 已提交
965 966
}

L
Liu Jicong 已提交
967
int32_t tmq_unsubscribe(tmq_t* tmq) {
L
Liu Jicong 已提交
968 969
  int32_t     rsp;
  int32_t     retryCnt = 0;
L
Liu Jicong 已提交
970
  tmq_list_t* lst = tmq_list_new();
L
Liu Jicong 已提交
971 972 973 974 975 976 977 978 979 980
  while (1) {
    rsp = tmq_subscribe(tmq, lst);
    if (rsp != TSDB_CODE_MND_CONSUMER_NOT_READY || retryCnt > 5) {
      break;
    } else {
      retryCnt++;
      taosMsleep(500);
    }
  }

L
Liu Jicong 已提交
981 982
  tmq_list_destroy(lst);
  return rsp;
X
Xiaoyu Wang 已提交
983 984
}

985 986 987 988 989 990
static void freeClientVgImpl(void* param) {
  SMqClientTopic* pTopic = param;
  taosMemoryFreeClear(pTopic->schema.pSchema);
  taosArrayDestroy(pTopic->vgs);
}

991
void tmqFreeImpl(void* handle) {
992 993
  tmq_t*  tmq = (tmq_t*)handle;
  int64_t id = tmq->consumerId;
L
Liu Jicong 已提交
994

995
  // TODO stop timer
L
Liu Jicong 已提交
996 997 998 999
  if (tmq->mqueue) {
    tmqClearUnhandleMsg(tmq);
    taosCloseQueue(tmq->mqueue);
  }
L
Liu Jicong 已提交
1000

H
Haojun Liao 已提交
1001 1002 1003 1004 1005
  if (tmq->delayedTask) {
    taosCloseQueue(tmq->delayedTask);
  }

  taosFreeQall(tmq->qall);
1006
  tsem_destroy(&tmq->rspSem);
L
Liu Jicong 已提交
1007

1008
  taosArrayDestroyEx(tmq->clientTopics, freeClientVgImpl);
1009 1010
  taos_close_internal(tmq->pTscObj);
  taosMemoryFree(tmq);
1011 1012

  tscDebug("consumer:0x%" PRIx64 " closed", id);
L
Liu Jicong 已提交
1013 1014
}

1015 1016 1017 1018 1019 1020 1021 1022 1023
static void tmqMgmtInit(void) {
  tmqInitRes = 0;
  tmqMgmt.timer = taosTmrInit(1000, 100, 360000, "TMQ");

  if (tmqMgmt.timer == NULL) {
    tmqInitRes = TSDB_CODE_OUT_OF_MEMORY;
  }

  tmqMgmt.rsetId = taosOpenRef(10000, tmqFreeImpl);
1024
  if (tmqMgmt.rsetId < 0) {
1025 1026 1027 1028
    tmqInitRes = terrno;
  }
}

L
Liu Jicong 已提交
1029
tmq_t* tmq_consumer_new(tmq_conf_t* conf, char* errstr, int32_t errstrLen) {
1030 1031 1032 1033
  taosThreadOnce(&tmqInit, tmqMgmtInit);
  if (tmqInitRes != 0) {
    terrno = tmqInitRes;
    return NULL;
L
Liu Jicong 已提交
1034 1035
  }

L
Liu Jicong 已提交
1036 1037
  tmq_t* pTmq = taosMemoryCalloc(1, sizeof(tmq_t));
  if (pTmq == NULL) {
L
Liu Jicong 已提交
1038
    terrno = TSDB_CODE_OUT_OF_MEMORY;
1039
    tscError("failed to create consumer, groupId:%s, code:%s", conf->groupId, terrstr());
L
Liu Jicong 已提交
1040 1041
    return NULL;
  }
L
Liu Jicong 已提交
1042

L
Liu Jicong 已提交
1043 1044 1045
  const char* user = conf->user == NULL ? TSDB_DEFAULT_USER : conf->user;
  const char* pass = conf->pass == NULL ? TSDB_DEFAULT_PASS : conf->pass;

L
Liu Jicong 已提交
1046 1047 1048
  pTmq->clientTopics = taosArrayInit(0, sizeof(SMqClientTopic));
  pTmq->mqueue = taosOpenQueue();
  pTmq->delayedTask = taosOpenQueue();
H
Haojun Liao 已提交
1049
  pTmq->qall = taosAllocateQall();
L
Liu Jicong 已提交
1050

X
Xiaoyu Wang 已提交
1051 1052
  if (pTmq->clientTopics == NULL || pTmq->mqueue == NULL || pTmq->qall == NULL || pTmq->delayedTask == NULL ||
      conf->groupId[0] == 0) {
L
Liu Jicong 已提交
1053
    terrno = TSDB_CODE_OUT_OF_MEMORY;
X
Xiaoyu Wang 已提交
1054
    tscError("consumer:0x%" PRIx64 " setup failed since %s, groupId:%s", pTmq->consumerId, terrstr(), pTmq->groupId);
1055
    goto _failed;
L
Liu Jicong 已提交
1056
  }
L
Liu Jicong 已提交
1057

L
Liu Jicong 已提交
1058 1059
  // init status
  pTmq->status = TMQ_CONSUMER_STATUS__INIT;
L
Liu Jicong 已提交
1060 1061
  pTmq->pollCnt = 0;
  pTmq->epoch = 0;
L
Liu Jicong 已提交
1062

L
Liu Jicong 已提交
1063 1064 1065
  // set conf
  strcpy(pTmq->clientId, conf->clientId);
  strcpy(pTmq->groupId, conf->groupId);
1066
  pTmq->withTbName = conf->withTbName;
L
Liu Jicong 已提交
1067
  pTmq->useSnapshot = conf->snapEnable;
L
Liu Jicong 已提交
1068
  pTmq->autoCommit = conf->autoCommit;
L
Liu Jicong 已提交
1069
  pTmq->autoCommitInterval = conf->autoCommitInterval;
L
Liu Jicong 已提交
1070 1071
  pTmq->commitCb = conf->commitCb;
  pTmq->commitCbUserParam = conf->commitCbUserParam;
L
Liu Jicong 已提交
1072 1073
  pTmq->resetOffsetCfg = conf->resetOffset;

1074 1075
  pTmq->hbBgEnable = conf->hbBgEnable;

L
Liu Jicong 已提交
1076
  // assign consumerId
L
Liu Jicong 已提交
1077
  pTmq->consumerId = tGenIdPI64();
X
Xiaoyu Wang 已提交
1078

L
Liu Jicong 已提交
1079 1080
  // init semaphore
  if (tsem_init(&pTmq->rspSem, 0, 0) != 0) {
1081
    tscError("consumer:0x %" PRIx64 " setup failed since %s, consumer group %s", pTmq->consumerId, terrstr(),
S
Shengliang Guan 已提交
1082
             pTmq->groupId);
1083
    goto _failed;
L
Liu Jicong 已提交
1084
  }
L
Liu Jicong 已提交
1085

L
Liu Jicong 已提交
1086 1087 1088
  // init connection
  pTmq->pTscObj = taos_connect_internal(conf->ip, user, pass, NULL, NULL, conf->port, CONN_TYPE__TMQ);
  if (pTmq->pTscObj == NULL) {
1089
    tscError("consumer:0x%" PRIx64 " setup failed since %s, groupId:%s", pTmq->consumerId, terrstr(), pTmq->groupId);
L
Liu Jicong 已提交
1090
    tsem_destroy(&pTmq->rspSem);
1091
    goto _failed;
L
Liu Jicong 已提交
1092
  }
L
Liu Jicong 已提交
1093

1094 1095
  pTmq->refId = taosAddRef(tmqMgmt.rsetId, pTmq);
  if (pTmq->refId < 0) {
1096
    goto _failed;
1097 1098
  }

1099
  if (pTmq->hbBgEnable) {
L
Liu Jicong 已提交
1100 1101
    int64_t* pRefId = taosMemoryMalloc(sizeof(int64_t));
    *pRefId = pTmq->refId;
1102
    pTmq->hbLiveTimer = taosTmrStart(tmqSendHbReq, 1000, pRefId, tmqMgmt.timer);
1103 1104
  }

1105 1106 1107
  char         buf[80] = {0};
  STqOffsetVal offset = {.type = pTmq->resetOffsetCfg};
  tFormatOffset(buf, tListLen(buf), &offset);
X
Xiaoyu Wang 已提交
1108 1109 1110 1111
  tscInfo("consumer:0x%" PRIx64 " is setup, refId:%" PRId64
          ", groupId:%s, snapshot:%d, autoCommit:%d, commitInterval:%dms, offset:%s, backgroudHB:%d",
          pTmq->consumerId, pTmq->refId, pTmq->groupId, pTmq->useSnapshot, pTmq->autoCommit, pTmq->autoCommitInterval,
          buf, pTmq->hbBgEnable);
L
Liu Jicong 已提交
1112

1113
  return pTmq;
1114

1115 1116
_failed:
  tmqFreeImpl(pTmq);
L
Liu Jicong 已提交
1117
  return NULL;
1118 1119
}

L
Liu Jicong 已提交
1120
int32_t tmq_subscribe(tmq_t* tmq, const tmq_list_t* topic_list) {
1121
  const int32_t   MAX_RETRY_COUNT = 120 * 4;  // let's wait for 4 mins at most
L
Liu Jicong 已提交
1122 1123 1124
  const SArray*   container = &topic_list->container;
  int32_t         sz = taosArrayGetSize(container);
  void*           buf = NULL;
L
Liu Jicong 已提交
1125
  SMsgSendInfo*   sendInfo = NULL;
L
Liu Jicong 已提交
1126
  SCMSubscribeReq req = {0};
1127
  int32_t         code = 0;
1128

1129
  tscDebug("consumer:0x%" PRIx64 " cgroup:%s, subscribe %d topics", tmq->consumerId, tmq->groupId, sz);
L
Liu Jicong 已提交
1130

1131
  req.consumerId = tmq->consumerId;
L
Liu Jicong 已提交
1132
  tstrncpy(req.clientId, tmq->clientId, 256);
L
Liu Jicong 已提交
1133
  tstrncpy(req.cgroup, tmq->groupId, TSDB_CGROUP_LEN);
1134 1135
  req.topicNames = taosArrayInit(sz, sizeof(void*));

1136 1137 1138 1139
  if (req.topicNames == NULL) {
    code = TSDB_CODE_OUT_OF_MEMORY;
    goto FAIL;
  }
1140

L
Liu Jicong 已提交
1141 1142
  for (int32_t i = 0; i < sz; i++) {
    char* topic = taosArrayGetP(container, i);
1143 1144

    SName name = {0};
L
Liu Jicong 已提交
1145 1146 1147 1148
    tNameSetDbName(&name, tmq->pTscObj->acctId, topic, strlen(topic));
    char* topicFName = taosMemoryCalloc(1, TSDB_TOPIC_FNAME_LEN);
    if (topicFName == NULL) {
      goto FAIL;
1149 1150
    }

1151
    tNameExtractFullName(&name, topicFName);
X
Xiaoyu Wang 已提交
1152
    tscDebug("consumer:0x%" PRIx64 " subscribe topic:%s", tmq->consumerId, topicFName);
L
Liu Jicong 已提交
1153 1154

    taosArrayPush(req.topicNames, &topicFName);
1155 1156
  }

L
Liu Jicong 已提交
1157
  int32_t tlen = tSerializeSCMSubscribeReq(NULL, &req);
1158

L
Liu Jicong 已提交
1159
  buf = taosMemoryMalloc(tlen);
1160 1161 1162 1163
  if (buf == NULL) {
    code = TSDB_CODE_OUT_OF_MEMORY;
    goto FAIL;
  }
L
Liu Jicong 已提交
1164

1165 1166 1167
  void* abuf = buf;
  tSerializeSCMSubscribeReq(&abuf, &req);

L
Liu Jicong 已提交
1168
  sendInfo = taosMemoryCalloc(1, sizeof(SMsgSendInfo));
1169 1170 1171 1172
  if (sendInfo == NULL) {
    code = TSDB_CODE_OUT_OF_MEMORY;
    goto FAIL;
  }
1173

1174
  SMqSubscribeCbParam param = { .rspErr = 0, .refId = tmq->refId, .epoch = tmq->epoch };
1175
  if (tsem_init(&param.rspSem, 0, 0) != 0) {
wmmhello's avatar
wmmhello 已提交
1176
    code = TSDB_CODE_TSC_INTERNAL_ERROR;
1177 1178
    goto FAIL;
  }
L
Liu Jicong 已提交
1179

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

L
Liu Jicong 已提交
1182 1183
  sendInfo->requestId = generateRequestId();
  sendInfo->requestObjRefId = 0;
L
Liu Jicong 已提交
1184 1185
  sendInfo->param = &param;
  sendInfo->fp = tmqSubscribeCb;
L
Liu Jicong 已提交
1186
  sendInfo->msgType = TDMT_MND_TMQ_SUBSCRIBE;
L
Liu Jicong 已提交
1187

1188 1189 1190 1191 1192
  SEpSet epSet = getEpSet_s(&tmq->pTscObj->pAppInfo->mgmtEp);

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

L
Liu Jicong 已提交
1193 1194
  // avoid double free if msg is sent
  buf = NULL;
L
Liu Jicong 已提交
1195
  sendInfo = NULL;
L
Liu Jicong 已提交
1196

L
Liu Jicong 已提交
1197 1198
  tsem_wait(&param.rspSem);
  tsem_destroy(&param.rspSem);
1199

1200 1201 1202 1203
  if (param.rspErr != 0) {
    code = param.rspErr;
    goto FAIL;
  }
L
Liu Jicong 已提交
1204

L
Liu Jicong 已提交
1205
  int32_t retryCnt = 0;
1206
  while (TSDB_CODE_MND_CONSUMER_NOT_READY == doAskEp(tmq)) {
1207
    if (retryCnt++ > MAX_RETRY_COUNT) {
1208
      tscError("consumer:0x%" PRIx64 ", mnd not ready for subscribe, max retry reached:%d", tmq->consumerId, retryCnt);
wmmhello's avatar
wmmhello 已提交
1209
      code = TSDB_CODE_TSC_INTERNAL_ERROR;
L
Liu Jicong 已提交
1210 1211
      goto FAIL;
    }
1212

X
Xiaoyu Wang 已提交
1213
    tscDebug("consumer:0x%" PRIx64 ", mnd not ready for subscribe, retry:%d in 500ms", tmq->consumerId, retryCnt);
L
Liu Jicong 已提交
1214 1215
    taosMsleep(500);
  }
1216

1217 1218
  // init ep timer
  if (tmq->epTimer == NULL) {
1219 1220 1221
    int64_t* pRefId1 = taosMemoryMalloc(sizeof(int64_t));
    *pRefId1 = tmq->refId;
    tmq->epTimer = taosTmrStart(tmqAssignAskEpTask, 1000, pRefId1, tmqMgmt.timer);
1222
  }
L
Liu Jicong 已提交
1223 1224

  // init auto commit timer
1225
  if (tmq->autoCommit && tmq->commitTimer == NULL) {
1226 1227 1228
    int64_t* pRefId2 = taosMemoryMalloc(sizeof(int64_t));
    *pRefId2 = tmq->refId;
    tmq->commitTimer = taosTmrStart(tmqAssignDelayedCommitTask, tmq->autoCommitInterval, pRefId2, tmqMgmt.timer);
L
Liu Jicong 已提交
1229 1230
  }

L
Liu Jicong 已提交
1231
FAIL:
L
Liu Jicong 已提交
1232
  taosArrayDestroyP(req.topicNames, taosMemoryFree);
L
Liu Jicong 已提交
1233
  taosMemoryFree(buf);
L
Liu Jicong 已提交
1234
  taosMemoryFree(sendInfo);
L
Liu Jicong 已提交
1235

L
Liu Jicong 已提交
1236
  return code;
1237 1238
}

L
Liu Jicong 已提交
1239
void tmq_conf_set_auto_commit_cb(tmq_conf_t* conf, tmq_commit_cb* cb, void* param) {
1240
  conf->commitCb = cb;
L
Liu Jicong 已提交
1241
  conf->commitCbUserParam = param;
L
Liu Jicong 已提交
1242
}
1243

1244
static int32_t tmqPollCb(void* param, SDataBuf* pMsg, int32_t code) {
X
Xiaoyu Wang 已提交
1245
  SMqPollCbParam* pParam = (SMqPollCbParam*)param;
1246 1247

  int64_t         refId = pParam->refId;
X
Xiaoyu Wang 已提交
1248
  SMqClientVg*    pVg = pParam->pVg;
L
Liu Jicong 已提交
1249
  SMqClientTopic* pTopic = pParam->pTopic;
1250

1251
  tmq_t* tmq = taosAcquireRef(tmqMgmt.rsetId, refId);
1252 1253 1254
  if (tmq == NULL) {
    taosMemoryFree(pParam);
    taosMemoryFree(pMsg->pData);
dengyihao's avatar
dengyihao 已提交
1255
    taosMemoryFree(pMsg->pEpSet);
1256 1257 1258 1259
    terrno = TSDB_CODE_TMQ_CONSUMER_CLOSED;
    return -1;
  }

H
Haojun Liao 已提交
1260 1261 1262 1263
  int32_t  epoch = pParam->epoch;
  int32_t  vgId = pParam->vgId;
  uint64_t requestId = pParam->requestId;

L
Liu Jicong 已提交
1264
  taosMemoryFree(pParam);
H
Haojun Liao 已提交
1265

L
Liu Jicong 已提交
1266
  if (code != 0) {
H
Haojun Liao 已提交
1267 1268
    tscWarn("consumer:0x%" PRIx64 " msg from vgId:%d discarded, epoch %d, since %s, reqId:0x%" PRIx64, tmq->consumerId,
            vgId, epoch, tstrerror(code), requestId);
H
Haojun Liao 已提交
1269

L
Liu Jicong 已提交
1270
    if (pMsg->pData) taosMemoryFree(pMsg->pData);
dengyihao's avatar
dengyihao 已提交
1271 1272
    if (pMsg->pEpSet) taosMemoryFree(pMsg->pEpSet);

H
Haojun Liao 已提交
1273
    // in case of consumer mismatch, wait for 500ms and retry
L
Liu Jicong 已提交
1274
    if (code == TSDB_CODE_TMQ_CONSUMER_MISMATCH) {
1275
      taosMsleep(500);
L
Liu Jicong 已提交
1276
      atomic_store_8(&tmq->status, TMQ_CONSUMER_STATUS__RECOVER);
X
Xiaoyu Wang 已提交
1277 1278
      tscDebug("consumer:0x%" PRIx64 " wait for the re-balance, wait for 500ms and set status to be RECOVER",
               tmq->consumerId);
H
Haojun Liao 已提交
1279
    } else if (code == TSDB_CODE_TQ_NO_COMMITTED_OFFSET) {
S
Shengliang Guan 已提交
1280
      SMqPollRspWrapper* pRspWrapper = taosAllocateQitem(sizeof(SMqPollRspWrapper), DEF_QITEM, 0);
L
Liu Jicong 已提交
1281
      if (pRspWrapper == NULL) {
H
Haojun Liao 已提交
1282 1283
        tscWarn("consumer:0x%" PRIx64 " msg from vgId:%d discarded, epoch %d since out of memory, reqId:0x%" PRIx64,
                tmq->consumerId, vgId, epoch, requestId);
L
Liu Jicong 已提交
1284 1285
        goto CREATE_MSG_FAIL;
      }
H
Haojun Liao 已提交
1286

L
Liu Jicong 已提交
1287 1288
      pRspWrapper->tmqRspType = TMQ_MSG_TYPE__END_RSP;
      taosWriteQitem(tmq->mqueue, pRspWrapper);
X
Xiaoyu Wang 已提交
1289
    } else if (code == TSDB_CODE_WAL_LOG_NOT_EXIST) {  // poll data while insert
1290
      taosMsleep(500);
L
Liu Jicong 已提交
1291
    }
H
Haojun Liao 已提交
1292

L
fix txn  
Liu Jicong 已提交
1293
    goto CREATE_MSG_FAIL;
L
Liu Jicong 已提交
1294 1295
  }

X
Xiaoyu Wang 已提交
1296
  int32_t msgEpoch = ((SMqRspHead*)pMsg->pData)->epoch;
1297 1298
  int32_t clientEpoch = atomic_load_32(&tmq->epoch);
  if (msgEpoch < clientEpoch) {
L
Liu Jicong 已提交
1299
    // do not write into queue since updating epoch reset
X
Xiaoyu Wang 已提交
1300 1301
    tscWarn("consumer:0x%" PRIx64
            " msg discard from vgId:%d since from earlier epoch, rsp epoch %d, current epoch %d, reqId:0x%" PRIx64,
1302
            tmq->consumerId, vgId, msgEpoch, clientEpoch, requestId);
H
Haojun Liao 已提交
1303

1304
    tsem_post(&tmq->rspSem);
1305 1306
    taosReleaseRef(tmqMgmt.rsetId, refId);

L
Liu Jicong 已提交
1307
    taosMemoryFree(pMsg->pData);
dengyihao's avatar
dengyihao 已提交
1308
    taosMemoryFree(pMsg->pEpSet);
X
Xiaoyu Wang 已提交
1309 1310 1311
    return 0;
  }

1312
  if (msgEpoch != clientEpoch) {
H
Haojun Liao 已提交
1313
    tscWarn("consumer:0x%" PRIx64 " mismatch rsp from vgId:%d, epoch %d, current epoch %d, reqId:0x%" PRIx64,
1314
            tmq->consumerId, vgId, msgEpoch, clientEpoch, requestId);
X
Xiaoyu Wang 已提交
1315 1316
  }

L
Liu Jicong 已提交
1317 1318 1319
  // handle meta rsp
  int8_t rspType = ((SMqRspHead*)pMsg->pData)->mqMsgType;

S
Shengliang Guan 已提交
1320
  SMqPollRspWrapper* pRspWrapper = taosAllocateQitem(sizeof(SMqPollRspWrapper), DEF_QITEM, 0);
L
Liu Jicong 已提交
1321
  if (pRspWrapper == NULL) {
L
Liu Jicong 已提交
1322
    taosMemoryFree(pMsg->pData);
dengyihao's avatar
dengyihao 已提交
1323
    taosMemoryFree(pMsg->pEpSet);
X
Xiaoyu Wang 已提交
1324 1325
    tscWarn("consumer:0x%" PRIx64 " msg discard from vgId:%d, epoch %d since out of memory", tmq->consumerId, vgId,
            epoch);
L
fix txn  
Liu Jicong 已提交
1326
    goto CREATE_MSG_FAIL;
L
Liu Jicong 已提交
1327
  }
L
Liu Jicong 已提交
1328

L
Liu Jicong 已提交
1329
  pRspWrapper->tmqRspType = rspType;
L
Liu Jicong 已提交
1330 1331
  pRspWrapper->vgHandle = pVg;
  pRspWrapper->topicHandle = pTopic;
H
Haojun Liao 已提交
1332
  pRspWrapper->reqId = requestId;
1333
  pRspWrapper->pEpset = pMsg->pEpSet;
1334
  pRspWrapper->vgId = pVg->vgId;
L
Liu Jicong 已提交
1335

1336
  pMsg->pEpSet = NULL;
L
Liu Jicong 已提交
1337
  if (rspType == TMQ_MSG_TYPE__POLL_RSP) {
L
Liu Jicong 已提交
1338 1339
    SDecoder decoder;
    tDecoderInit(&decoder, POINTER_SHIFT(pMsg->pData, sizeof(SMqRspHead)), pMsg->len - sizeof(SMqRspHead));
1340
    tDecodeMqDataRsp(&decoder, &pRspWrapper->dataRsp);
wmmhello's avatar
wmmhello 已提交
1341
    tDecoderClear(&decoder);
L
Liu Jicong 已提交
1342
    memcpy(&pRspWrapper->dataRsp, pMsg->pData, sizeof(SMqRspHead));
L
Liu Jicong 已提交
1343

H
Haojun Liao 已提交
1344 1345
    char buf[80];
    tFormatOffset(buf, 80, &pRspWrapper->dataRsp.rspOffset);
H
Haojun Liao 已提交
1346
    tscDebug("consumer:0x%" PRIx64 " recv poll rsp, vgId:%d, req ver:%" PRId64 ", rsp:%s type %d, reqId:0x%" PRIx64,
H
Haojun Liao 已提交
1347
             tmq->consumerId, vgId, pRspWrapper->dataRsp.reqOffset.version, buf, rspType, requestId);
L
Liu Jicong 已提交
1348
  } else if (rspType == TMQ_MSG_TYPE__POLL_META_RSP) {
1349 1350
    SDecoder decoder;
    tDecoderInit(&decoder, POINTER_SHIFT(pMsg->pData, sizeof(SMqRspHead)), pMsg->len - sizeof(SMqRspHead));
1351
    tDecodeMqMetaRsp(&decoder, &pRspWrapper->metaRsp);
1352
    tDecoderClear(&decoder);
L
Liu Jicong 已提交
1353
    memcpy(&pRspWrapper->metaRsp, pMsg->pData, sizeof(SMqRspHead));
L
Liu Jicong 已提交
1354 1355 1356 1357 1358 1359
  } else if (rspType == TMQ_MSG_TYPE__TAOSX_RSP) {
    SDecoder decoder;
    tDecoderInit(&decoder, POINTER_SHIFT(pMsg->pData, sizeof(SMqRspHead)), pMsg->len - sizeof(SMqRspHead));
    tDecodeSTaosxRsp(&decoder, &pRspWrapper->taosxRsp);
    tDecoderClear(&decoder);
    memcpy(&pRspWrapper->taosxRsp, pMsg->pData, sizeof(SMqRspHead));
X
Xiaoyu Wang 已提交
1360 1361
  } else {  // invalid rspType
    tscError("consumer:0x%" PRIx64 " invalid rsp msg received, type:%d ignored", tmq->consumerId, rspType);
L
Liu Jicong 已提交
1362
  }
L
Liu Jicong 已提交
1363

L
Liu Jicong 已提交
1364
  taosMemoryFree(pMsg->pData);
H
Haojun Liao 已提交
1365
  taosWriteQitem(tmq->mqueue, pRspWrapper);
L
Liu Jicong 已提交
1366

1367
  int32_t total = taosQueueItemSize(tmq->mqueue);
H
Haojun Liao 已提交
1368
  tscDebug("consumer:0x%" PRIx64 " put poll res into mqueue, type:%d, vgId:%d, total in queue:%d, reqId:0x%" PRIx64,
1369
           tmq->consumerId, rspType, vgId, total, requestId);
H
Haojun Liao 已提交
1370

1371
  tsem_post(&tmq->rspSem);
1372 1373
  taosReleaseRef(tmqMgmt.rsetId, refId);

L
Liu Jicong 已提交
1374
  return 0;
H
Haojun Liao 已提交
1375

L
fix txn  
Liu Jicong 已提交
1376
CREATE_MSG_FAIL:
L
Liu Jicong 已提交
1377
  if (epoch == tmq->epoch) {
L
Liu Jicong 已提交
1378 1379
    atomic_store_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE);
  }
H
Haojun Liao 已提交
1380

1381
  tsem_post(&tmq->rspSem);
1382 1383
  taosReleaseRef(tmqMgmt.rsetId, refId);

L
Liu Jicong 已提交
1384
  return -1;
1385 1386
}

H
Haojun Liao 已提交
1387 1388 1389 1390 1391
typedef struct SVgroupSaveInfo {
  STqOffsetVal offset;
  int64_t      numOfRows;
} SVgroupSaveInfo;

H
Haojun Liao 已提交
1392 1393 1394 1395 1396 1397
static void initClientTopicFromRsp(SMqClientTopic* pTopic, SMqSubTopicEp* pTopicEp, SHashObj* pVgOffsetHashMap,
                                   tmq_t* tmq) {
  pTopic->schema = pTopicEp->schema;
  pTopicEp->schema.nCols = 0;
  pTopicEp->schema.pSchema = NULL;

X
Xiaoyu Wang 已提交
1398
  char    vgKey[TSDB_TOPIC_FNAME_LEN + 22];
H
Haojun Liao 已提交
1399 1400 1401 1402 1403 1404 1405 1406 1407 1408
  int32_t vgNumGet = taosArrayGetSize(pTopicEp->vgs);

  tstrncpy(pTopic->topicName, pTopicEp->topic, TSDB_TOPIC_FNAME_LEN);
  tstrncpy(pTopic->db, pTopicEp->db, TSDB_DB_FNAME_LEN);

  tscDebug("consumer:0x%" PRIx64 ", update topic:%s, numOfVgs:%d", tmq->consumerId, pTopic->topicName, vgNumGet);
  pTopic->vgs = taosArrayInit(vgNumGet, sizeof(SMqClientVg));

  for (int32_t j = 0; j < vgNumGet; j++) {
    SMqSubVgEp* pVgEp = taosArrayGet(pTopicEp->vgs, j);
H
Haojun Liao 已提交
1409 1410

    makeTopicVgroupKey(vgKey, pTopic->topicName, pVgEp->vgId);
H
Haojun Liao 已提交
1411
    SVgroupSaveInfo* pInfo = taosHashGet(pVgOffsetHashMap, vgKey, strlen(vgKey));
H
Haojun Liao 已提交
1412

X
Xiaoyu Wang 已提交
1413 1414
    int64_t      numOfRows = 0;
    STqOffsetVal offsetNew = {.type = tmq->resetOffsetCfg};
H
Haojun Liao 已提交
1415 1416 1417
    if (pInfo != NULL) {
      offsetNew = pInfo->offset;
      numOfRows = pInfo->numOfRows;
H
Haojun Liao 已提交
1418 1419 1420 1421 1422 1423 1424 1425
    }

    SMqClientVg clientVg = {
        .pollCnt = 0,
        .vgId = pVgEp->vgId,
        .epSet = pVgEp->epSet,
        .vgStatus = TMQ_VG_STATUS__IDLE,
        .vgSkipCnt = 0,
H
Haojun Liao 已提交
1426
        .emptyBlockReceiveTs = 0,
H
Haojun Liao 已提交
1427
        .numOfRows = numOfRows,
H
Haojun Liao 已提交
1428 1429
    };

H
Haojun Liao 已提交
1430 1431 1432 1433
    clientVg.offsetInfo.currentOffset = offsetNew;
    clientVg.offsetInfo.committedOffset = offsetNew;
    clientVg.offsetInfo.walVerBegin = -1;
    clientVg.offsetInfo.walVerEnd = -1;
H
Haojun Liao 已提交
1434 1435 1436 1437 1438 1439 1440 1441 1442 1443 1444 1445 1446
    taosArrayPush(pTopic->vgs, &clientVg);
  }
}

static void freeClientVgInfo(void* param) {
  SMqClientTopic* pTopic = param;
  if (pTopic->schema.nCols) {
    taosMemoryFreeClear(pTopic->schema.pSchema);
  }

  taosArrayDestroy(pTopic->vgs);
}

1447
static bool doUpdateLocalEp(tmq_t* tmq, int32_t epoch, const SMqAskEpRsp* pRsp) {
1448 1449
  bool set = false;

1450
  int32_t topicNumCur = taosArrayGetSize(tmq->clientTopics);
1451
  int32_t topicNumGet = taosArrayGetSize(pRsp->topics);
1452

X
Xiaoyu Wang 已提交
1453 1454
  char vgKey[TSDB_TOPIC_FNAME_LEN + 22];
  tscDebug("consumer:0x%" PRIx64 " update ep epoch from %d to epoch %d, incoming topics:%d, existed topics:%d",
1455
           tmq->consumerId, tmq->epoch, epoch, topicNumGet, topicNumCur);
wmmhello's avatar
wmmhello 已提交
1456 1457 1458
  if (epoch <= tmq->epoch) {
    return false;
  }
1459 1460 1461 1462 1463 1464

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

H
Haojun Liao 已提交
1465 1466
  SHashObj* pVgOffsetHashMap = taosHashInit(64, MurmurHash3_32, false, HASH_NO_LOCK);
  if (pVgOffsetHashMap == NULL) {
1467 1468 1469
    taosArrayDestroy(newTopics);
    return false;
  }
1470

H
Haojun Liao 已提交
1471
  // todo extract method
1472 1473 1474 1475 1476
  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);
1477
      tscDebug("consumer:0x%" PRIx64 ", new vg num: %d", tmq->consumerId, vgNumCur);
1478 1479
      for (int32_t j = 0; j < vgNumCur; j++) {
        SMqClientVg* pVgCur = taosArrayGet(pTopicCur->vgs, j);
H
Haojun Liao 已提交
1480 1481
        makeTopicVgroupKey(vgKey, pTopicCur->topicName, pVgCur->vgId);

L
Liu Jicong 已提交
1482
        char buf[80];
H
Haojun Liao 已提交
1483
        tFormatOffset(buf, 80, &pVgCur->offsetInfo.currentOffset);
X
Xiaoyu Wang 已提交
1484 1485
        tscDebug("consumer:0x%" PRIx64 ", epoch:%d vgId:%d vgKey:%s, offset:%s", tmq->consumerId, epoch, pVgCur->vgId,
                 vgKey, buf);
H
Haojun Liao 已提交
1486

H
Haojun Liao 已提交
1487
        SVgroupSaveInfo info = {.offset = pVgCur->offsetInfo.currentOffset, .numOfRows = pVgCur->numOfRows};
H
Haojun Liao 已提交
1488
        taosHashPut(pVgOffsetHashMap, vgKey, strlen(vgKey), &info, sizeof(SVgroupSaveInfo));
1489 1490 1491 1492 1493 1494 1495
      }
    }
  }

  for (int32_t i = 0; i < topicNumGet; i++) {
    SMqClientTopic topic = {0};
    SMqSubTopicEp* pTopicEp = taosArrayGet(pRsp->topics, i);
H
Haojun Liao 已提交
1496
    initClientTopicFromRsp(&topic, pTopicEp, pVgOffsetHashMap, tmq);
1497 1498
    taosArrayPush(newTopics, &topic);
  }
1499

H
Haojun Liao 已提交
1500 1501
  taosHashCleanup(pVgOffsetHashMap);

1502
  // destroy current buffered existed topics info
1503
  if (tmq->clientTopics) {
H
Haojun Liao 已提交
1504
    taosArrayDestroyEx(tmq->clientTopics, freeClientVgInfo);
X
Xiaoyu Wang 已提交
1505
  }
H
Haojun Liao 已提交
1506
  tmq->clientTopics = newTopics;
1507

X
Xiaoyu Wang 已提交
1508
  int8_t flag = (topicNumGet == 0) ? TMQ_CONSUMER_STATUS__NO_TOPIC : TMQ_CONSUMER_STATUS__READY;
H
Haojun Liao 已提交
1509
  atomic_store_8(&tmq->status, flag);
X
Xiaoyu Wang 已提交
1510
  atomic_store_32(&tmq->epoch, epoch);
H
Haojun Liao 已提交
1511

1512
  tscDebug("consumer:0x%" PRIx64 " update topic info completed", tmq->consumerId);
X
Xiaoyu Wang 已提交
1513 1514 1515
  return set;
}

1516
int32_t askEpCallbackFn(void* param, SDataBuf* pMsg, int32_t code) {
1517
  SMqAskEpCbParam* pParam = (SMqAskEpCbParam*)param;
1518 1519 1520
  tmq_t*           tmq = taosAcquireRef(tmqMgmt.rsetId, pParam->refId);

  if (tmq == NULL) {
1521 1522 1523
    terrno = TSDB_CODE_TMQ_CONSUMER_CLOSED;
    pParam->pUserFn(tmq, terrno, NULL, pParam->pParam);

1524
    taosMemoryFree(pMsg->pData);
dengyihao's avatar
dengyihao 已提交
1525
    taosMemoryFree(pMsg->pEpSet);
1526 1527
    taosMemoryFree(pParam);
    return terrno;
1528 1529
  }

H
Haojun Liao 已提交
1530
  if (code != TSDB_CODE_SUCCESS) {
1531 1532 1533 1534 1535 1536 1537 1538 1539
    tscError("consumer:0x%" PRIx64 ", get topic endpoint error, code:%s", tmq->consumerId, tstrerror(code));
    pParam->pUserFn(tmq, code, NULL, pParam->pParam);

    taosReleaseRef(tmqMgmt.rsetId, pParam->refId);

    taosMemoryFree(pMsg->pData);
    taosMemoryFree(pMsg->pEpSet);
    taosMemoryFree(pParam);
    return code;
1540
  }
L
Liu Jicong 已提交
1541

L
Liu Jicong 已提交
1542
  // tmq's epoch is monotonically increase,
L
Liu Jicong 已提交
1543
  // so it's safe to discard any old epoch msg.
L
Liu Jicong 已提交
1544
  // Epoch will only increase when received newer epoch ep msg
L
Liu Jicong 已提交
1545 1546 1547
  SMqRspHead* head = pMsg->pData;
  int32_t     epoch = atomic_load_32(&tmq->epoch);
  if (head->epoch <= epoch) {
1548 1549
    tscDebug("consumer:0x%" PRIx64 ", recv ep, msg epoch %d, current epoch %d, no need to update local ep",
             tmq->consumerId, head->epoch, epoch);
1550

1551 1552 1553 1554 1555 1556 1557 1558
    if (tmq->status == TMQ_CONSUMER_STATUS__RECOVER) {
      SMqAskEpRsp rsp;
      tDecodeSMqAskEpRsp(POINTER_SHIFT(pMsg->pData, sizeof(SMqRspHead)), &rsp);
      int8_t flag = (taosArrayGetSize(rsp.topics) == 0) ? TMQ_CONSUMER_STATUS__NO_TOPIC : TMQ_CONSUMER_STATUS__READY;
      atomic_store_8(&tmq->status, flag);
      tDeleteSMqAskEpRsp(&rsp);
    }

X
Xiaoyu Wang 已提交
1559
  } else {
1560 1561
    tscDebug("consumer:0x%" PRIx64 ", recv ep, msg epoch %d, current epoch %d, update local ep", tmq->consumerId,
             head->epoch, epoch);
1562
  }
L
Liu Jicong 已提交
1563

1564
  pParam->pUserFn(tmq, code, pMsg, pParam->pParam);
1565 1566
  taosReleaseRef(tmqMgmt.rsetId, pParam->refId);

dengyihao's avatar
dengyihao 已提交
1567
  taosMemoryFree(pMsg->pEpSet);
L
Liu Jicong 已提交
1568
  taosMemoryFree(pMsg->pData);
1569
  taosMemoryFree(pParam);
L
Liu Jicong 已提交
1570
  return code;
1571 1572
}

L
Liu Jicong 已提交
1573
void tmqBuildConsumeReqImpl(SMqPollReq* pReq, tmq_t* tmq, int64_t timeout, SMqClientTopic* pTopic, SMqClientVg* pVg) {
L
Liu Jicong 已提交
1574 1575 1576 1577
  int32_t groupLen = strlen(tmq->groupId);
  memcpy(pReq->subKey, tmq->groupId, groupLen);
  pReq->subKey[groupLen] = TMQ_SEPARATOR;
  strcpy(pReq->subKey + groupLen + 1, pTopic->topicName);
1578

1579
  pReq->withTbName = tmq->withTbName;
L
Liu Jicong 已提交
1580
  pReq->consumerId = tmq->consumerId;
1581
  pReq->timeout = timeout;
X
Xiaoyu Wang 已提交
1582
  pReq->epoch = tmq->epoch;
H
Haojun Liao 已提交
1583
  pReq->reqOffset = pVg->offsetInfo.currentOffset;
D
dapan1121 已提交
1584
  pReq->head.vgId = pVg->vgId;
1585 1586
  pReq->useSnapshot = tmq->useSnapshot;
  pReq->reqId = generateRequestId();
1587 1588
}

L
Liu Jicong 已提交
1589 1590
SMqMetaRspObj* tmqBuildMetaRspFromWrapper(SMqPollRspWrapper* pWrapper) {
  SMqMetaRspObj* pRspObj = taosMemoryCalloc(1, sizeof(SMqMetaRspObj));
L
Liu Jicong 已提交
1591
  pRspObj->resType = RES_TYPE__TMQ_META;
L
Liu Jicong 已提交
1592 1593 1594 1595 1596 1597 1598 1599
  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;
}

1600
SMqRspObj* tmqBuildRspFromWrapper(SMqPollRspWrapper* pWrapper, SMqClientVg* pVg, int64_t* numOfRows) {
L
Liu Jicong 已提交
1601 1602
  SMqRspObj* pRspObj = taosMemoryCalloc(1, sizeof(SMqRspObj));
  pRspObj->resType = RES_TYPE__TMQ;
1603

1604
  (*numOfRows) = 0;
L
Liu Jicong 已提交
1605 1606
  tstrncpy(pRspObj->topic, pWrapper->topicHandle->topicName, TSDB_TOPIC_FNAME_LEN);
  tstrncpy(pRspObj->db, pWrapper->topicHandle->db, TSDB_DB_FNAME_LEN);
1607

L
Liu Jicong 已提交
1608
  pRspObj->vgId = pWrapper->vgHandle->vgId;
L
Liu Jicong 已提交
1609
  pRspObj->resIter = -1;
L
Liu Jicong 已提交
1610
  memcpy(&pRspObj->rsp, &pWrapper->dataRsp, sizeof(SMqDataRsp));
L
Liu Jicong 已提交
1611

L
Liu Jicong 已提交
1612 1613
  pRspObj->resInfo.totalRows = 0;
  pRspObj->resInfo.precision = TSDB_TIME_PRECISION_MILLI;
1614

L
Liu Jicong 已提交
1615
  if (!pWrapper->dataRsp.withSchema) {
L
Liu Jicong 已提交
1616 1617
    setResSchemaInfo(&pRspObj->resInfo, pWrapper->topicHandle->schema.pSchema, pWrapper->topicHandle->schema.nCols);
  }
L
Liu Jicong 已提交
1618

1619
  // extract the rows in this data packet
X
Xiaoyu Wang 已提交
1620
  for (int32_t i = 0; i < pRspObj->rsp.blockNum; ++i) {
1621
    SRetrieveTableRsp* pRetrieve = (SRetrieveTableRsp*)taosArrayGetP(pRspObj->rsp.blockData, i);
X
Xiaoyu Wang 已提交
1622
    int64_t            rows = htobe64(pRetrieve->numOfRows);
1623
    pVg->numOfRows += rows;
1624
    (*numOfRows) += rows;
1625 1626
  }

L
Liu Jicong 已提交
1627
  return pRspObj;
X
Xiaoyu Wang 已提交
1628 1629
}

L
Liu Jicong 已提交
1630 1631
SMqTaosxRspObj* tmqBuildTaosxRspFromWrapper(SMqPollRspWrapper* pWrapper) {
  SMqTaosxRspObj* pRspObj = taosMemoryCalloc(1, sizeof(SMqTaosxRspObj));
1632
  pRspObj->resType = RES_TYPE__TMQ_METADATA;
L
Liu Jicong 已提交
1633 1634 1635 1636
  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;
  pRspObj->resIter = -1;
1637
  memcpy(&pRspObj->rsp, &pWrapper->taosxRsp, sizeof(STaosxRsp));
L
Liu Jicong 已提交
1638 1639 1640

  pRspObj->resInfo.totalRows = 0;
  pRspObj->resInfo.precision = TSDB_TIME_PRECISION_MILLI;
1641
  if (!pWrapper->taosxRsp.withSchema) {
L
Liu Jicong 已提交
1642 1643 1644 1645 1646 1647
    setResSchemaInfo(&pRspObj->resInfo, pWrapper->topicHandle->schema.pSchema, pWrapper->topicHandle->schema.nCols);
  }

  return pRspObj;
}

1648 1649 1650 1651 1652 1653 1654 1655 1656 1657 1658 1659 1660 1661 1662 1663 1664 1665 1666 1667 1668 1669 1670 1671 1672 1673 1674 1675 1676 1677 1678 1679 1680
static int32_t handleErrorBeforePoll(SMqClientVg* pVg, tmq_t* pTmq) {
  atomic_store_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE);
  tsem_post(&pTmq->rspSem);
  return -1;
}

static int32_t doTmqPollImpl(tmq_t* pTmq, SMqClientTopic* pTopic, SMqClientVg* pVg, int64_t timeout) {
  SMqPollReq req = {0};
  tmqBuildConsumeReqImpl(&req, pTmq, timeout, pTopic, pVg);

  int32_t msgSize = tSerializeSMqPollReq(NULL, 0, &req);
  if (msgSize < 0) {
    return handleErrorBeforePoll(pVg, pTmq);
  }

  char* msg = taosMemoryCalloc(1, msgSize);
  if (NULL == msg) {
    return handleErrorBeforePoll(pVg, pTmq);
  }

  if (tSerializeSMqPollReq(msg, msgSize, &req) < 0) {
    taosMemoryFree(msg);
    return handleErrorBeforePoll(pVg, pTmq);
  }

  SMqPollCbParam* pParam = taosMemoryMalloc(sizeof(SMqPollCbParam));
  if (pParam == NULL) {
    taosMemoryFree(msg);
    return handleErrorBeforePoll(pVg, pTmq);
  }

  pParam->refId = pTmq->refId;
  pParam->epoch = pTmq->epoch;
H
Haojun Liao 已提交
1681
  pParam->pVg = pVg;
1682 1683
  pParam->pTopic = pTopic;
  pParam->vgId = pVg->vgId;
H
Haojun Liao 已提交
1684
  pParam->requestId = req.reqId;
1685 1686 1687 1688 1689 1690 1691 1692

  SMsgSendInfo* sendInfo = taosMemoryCalloc(1, sizeof(SMsgSendInfo));
  if (sendInfo == NULL) {
    taosMemoryFree(pParam);
    taosMemoryFree(msg);
    return handleErrorBeforePoll(pVg, pTmq);
  }

H
Haojun Liao 已提交
1693
  sendInfo->msgInfo = (SDataBuf){.pData = msg, .len = msgSize, .handle = NULL};
1694 1695 1696 1697 1698 1699 1700 1701
  sendInfo->requestId = req.reqId;
  sendInfo->requestObjRefId = 0;
  sendInfo->param = pParam;
  sendInfo->fp = tmqPollCb;
  sendInfo->msgType = TDMT_VND_TMQ_CONSUME;

  int64_t transporterId = 0;
  char    offsetFormatBuf[80];
H
Haojun Liao 已提交
1702
  tFormatOffset(offsetFormatBuf, tListLen(offsetFormatBuf), &pVg->offsetInfo.currentOffset);
1703

X
Xiaoyu Wang 已提交
1704 1705
  tscDebug("consumer:0x%" PRIx64 " send poll to %s vgId:%d, epoch %d, req:%s, reqId:0x%" PRIx64, pTmq->consumerId,
           pTopic->topicName, pVg->vgId, pTmq->epoch, offsetFormatBuf, req.reqId);
1706 1707 1708 1709 1710 1711 1712 1713
  asyncSendMsgToServer(pTmq->pTscObj->pAppInfo->pTransporter, &pVg->epSet, &transporterId, sendInfo);

  pVg->pollCnt++;
  pTmq->pollCnt++;

  return TSDB_CODE_SUCCESS;
}

1714
// broadcast the poll request to all related vnodes
H
Haojun Liao 已提交
1715
static int32_t tmqPollImpl(tmq_t* tmq, int64_t timeout) {
1716
  int32_t numOfTopics = taosArrayGetSize(tmq->clientTopics);
X
Xiaoyu Wang 已提交
1717
  tscDebug("consumer:0x%" PRIx64 " start to poll data, numOfTopics:%d", tmq->consumerId, numOfTopics);
1718 1719

  for (int i = 0; i < numOfTopics; i++) {
X
Xiaoyu Wang 已提交
1720
    SMqClientTopic* pTopic = taosArrayGet(tmq->clientTopics, i);
X
Xiaoyu Wang 已提交
1721
    int32_t         numOfVg = taosArrayGetSize(pTopic->vgs);
1722 1723

    for (int j = 0; j < numOfVg; j++) {
X
Xiaoyu Wang 已提交
1724
      SMqClientVg* pVg = taosArrayGet(pTopic->vgs, j);
X
Xiaoyu Wang 已提交
1725 1726 1727
      if (taosGetTimestampMs() - pVg->emptyBlockReceiveTs < EMPTY_BLOCK_POLL_IDLE_DURATION) {  // less than 100ms
        tscTrace("consumer:0x%" PRIx64 " epoch %d, vgId:%d idle for 10ms before start next poll", tmq->consumerId,
                 tmq->epoch, pVg->vgId);
H
Haojun Liao 已提交
1728 1729 1730
        continue;
      }

1731
      int32_t vgStatus = atomic_val_compare_exchange_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE, TMQ_VG_STATUS__WAIT);
1732
      if (vgStatus == TMQ_VG_STATUS__WAIT) {
L
Liu Jicong 已提交
1733
        int32_t vgSkipCnt = atomic_add_fetch_32(&pVg->vgSkipCnt, 1);
H
Haojun Liao 已提交
1734
        tscTrace("consumer:0x%" PRIx64 " epoch %d wait poll-rsp, skip vgId:%d skip cnt %d", tmq->consumerId, tmq->epoch,
X
Xiaoyu Wang 已提交
1735
                 pVg->vgId, vgSkipCnt);
X
Xiaoyu Wang 已提交
1736 1737
        continue;
      }
1738

L
Liu Jicong 已提交
1739
      atomic_store_32(&pVg->vgSkipCnt, 0);
1740 1741 1742
      int32_t code = doTmqPollImpl(tmq, pTopic, pVg, timeout);
      if (code != TSDB_CODE_SUCCESS) {
        return code;
D
dapan1121 已提交
1743
      }
X
Xiaoyu Wang 已提交
1744 1745
    }
  }
1746

1747
  tscDebug("consumer:0x%" PRIx64 " end to poll data", tmq->consumerId);
X
Xiaoyu Wang 已提交
1748 1749 1750
  return 0;
}

H
Haojun Liao 已提交
1751
static int32_t tmqHandleNoPollRsp(tmq_t* tmq, SMqRspWrapper* rspWrapper, bool* pReset) {
L
Liu Jicong 已提交
1752
  if (rspWrapper->tmqRspType == TMQ_MSG_TYPE__EP_RSP) {
L
fix  
Liu Jicong 已提交
1753
    /*printf("ep %d %d\n", rspMsg->head.epoch, tmq->epoch);*/
L
Liu Jicong 已提交
1754 1755
    if (rspWrapper->epoch > atomic_load_32(&tmq->epoch)) {
      SMqAskEpRspWrapper* pEpRspWrapper = (SMqAskEpRspWrapper*)rspWrapper;
L
Liu Jicong 已提交
1756
      SMqAskEpRsp*        rspMsg = &pEpRspWrapper->msg;
1757
      doUpdateLocalEp(tmq, rspWrapper->epoch, rspMsg);
L
temp  
Liu Jicong 已提交
1758
      /*tmqClearUnhandleMsg(tmq);*/
L
Liu Jicong 已提交
1759
      tDeleteSMqAskEpRsp(rspMsg);
X
Xiaoyu Wang 已提交
1760 1761
      *pReset = true;
    } else {
L
Liu Jicong 已提交
1762
      tmqFreeRspWrapper(rspWrapper);
X
Xiaoyu Wang 已提交
1763 1764 1765 1766 1767 1768 1769 1770
      *pReset = false;
    }
  } else {
    return -1;
  }
  return 0;
}

H
Haojun Liao 已提交
1771
static void* tmqHandleAllRsp(tmq_t* tmq, int64_t timeout, bool pollIfReset) {
H
Haojun Liao 已提交
1772
  tscDebug("consumer:0x%" PRIx64 " start to handle the rsp, total:%d", tmq->consumerId, tmq->qall->numOfItems);
1773

X
Xiaoyu Wang 已提交
1774
  while (1) {
1775 1776
    SMqRspWrapper* pRspWrapper = NULL;
    taosGetQitem(tmq->qall, (void**)&pRspWrapper);
1777

1778
    if (pRspWrapper == NULL) {
L
Liu Jicong 已提交
1779
      taosReadAllQitems(tmq->mqueue, tmq->qall);
1780 1781
      taosGetQitem(tmq->qall, (void**)&pRspWrapper);
      if (pRspWrapper == NULL) {
L
Liu Jicong 已提交
1782 1783
        return NULL;
      }
X
Xiaoyu Wang 已提交
1784 1785
    }

X
Xiaoyu Wang 已提交
1786
    tscDebug("consumer:0x%" PRIx64 " handle rsp, type:%d", tmq->consumerId, pRspWrapper->tmqRspType);
H
Haojun Liao 已提交
1787

1788 1789
    if (pRspWrapper->tmqRspType == TMQ_MSG_TYPE__END_RSP) {
      taosFreeQitem(pRspWrapper);
L
Liu Jicong 已提交
1790
      terrno = TSDB_CODE_TQ_NO_COMMITTED_OFFSET;
H
Haojun Liao 已提交
1791
      tscError("consumer:0x%" PRIx64 " unexpected rsp from poll, code:%s", tmq->consumerId, tstrerror(terrno));
L
Liu Jicong 已提交
1792
      return NULL;
1793 1794
    } else if (pRspWrapper->tmqRspType == TMQ_MSG_TYPE__POLL_RSP) {
      SMqPollRspWrapper* pollRspWrapper = (SMqPollRspWrapper*)pRspWrapper;
H
Haojun Liao 已提交
1795

X
Xiaoyu Wang 已提交
1796
      int32_t     consumerEpoch = atomic_load_32(&tmq->epoch);
1797 1798 1799
      SMqDataRsp* pDataRsp = &pollRspWrapper->dataRsp;

      if (pDataRsp->head.epoch == consumerEpoch) {
L
Liu Jicong 已提交
1800
        SMqClientVg* pVg = pollRspWrapper->vgHandle;
1801 1802 1803 1804 1805

        // update the epset
        if (pollRspWrapper->pEpset != NULL) {
          SEp* pEp = GET_ACTIVE_EP(pollRspWrapper->pEpset);
          SEp* pOld = GET_ACTIVE_EP(&(pVg->epSet));
X
Xiaoyu Wang 已提交
1806 1807
          tscDebug("consumer:0x%" PRIx64 " update epset vgId:%d, ep:%s:%d, old ep:%s:%d", tmq->consumerId, pVg->vgId,
                   pEp->fqdn, pEp->port, pOld->fqdn, pOld->port);
1808 1809 1810
          pVg->epSet = *pollRspWrapper->pEpset;
        }

1811
        // update the local offset value only for the returned values.
H
Haojun Liao 已提交
1812
        pVg->offsetInfo.currentOffset = pDataRsp->rspOffset;
1813 1814

        // update the status
X
Xiaoyu Wang 已提交
1815
        atomic_store_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE);
H
Haojun Liao 已提交
1816

1817 1818 1819
        // update the valid wal version range
        pVg->offsetInfo.walVerBegin = pDataRsp->head.walsver;
        pVg->offsetInfo.walVerEnd = pDataRsp->head.walever;
1820
        pVg->receiveInfo = true;
1821

1822 1823 1824
        char buf[80];
        tFormatOffset(buf, 80, &pDataRsp->rspOffset);
        if (pDataRsp->blockNum == 0) {
X
Xiaoyu Wang 已提交
1825 1826 1827
          tscDebug("consumer:0x%" PRIx64 " empty block received, vgId:%d, offset:%s, vg total:%" PRId64
                   " total:%" PRId64 " reqId:0x%" PRIx64,
                   tmq->consumerId, pVg->vgId, buf, pVg->numOfRows, tmq->totalRows, pollRspWrapper->reqId);
1828
          pRspWrapper = tmqFreeRspWrapper(pRspWrapper);
L
Liu Jicong 已提交
1829
          taosFreeQitem(pollRspWrapper);
1830
        } else {  // build rsp
X
Xiaoyu Wang 已提交
1831
          int64_t    numOfRows = 0;
1832
          SMqRspObj* pRsp = tmqBuildRspFromWrapper(pollRspWrapper, pVg, &numOfRows);
H
Haojun Liao 已提交
1833 1834
          tmq->totalRows += numOfRows;

H
Haojun Liao 已提交
1835
          tscDebug("consumer:0x%" PRIx64 " process poll rsp, vgId:%d, offset:%s, blocks:%d, rows:%" PRId64
H
Haojun Liao 已提交
1836
                   " vg total:%" PRId64 " total:%" PRId64 ", reqId:0x%" PRIx64,
H
Haojun Liao 已提交
1837
                   tmq->consumerId, pVg->vgId, buf, pDataRsp->blockNum, numOfRows, pVg->numOfRows, tmq->totalRows,
H
Haojun Liao 已提交
1838
                   pollRspWrapper->reqId);
1839 1840 1841
          taosFreeQitem(pollRspWrapper);
          return pRsp;
        }
X
Xiaoyu Wang 已提交
1842
      } else {
H
Haojun Liao 已提交
1843
        tscDebug("consumer:0x%" PRIx64 " vgId:%d msg discard since epoch mismatch: msg epoch %d, consumer epoch %d",
1844
                 tmq->consumerId, pollRspWrapper->vgId, pDataRsp->head.epoch, consumerEpoch);
1845
        pRspWrapper = tmqFreeRspWrapper(pRspWrapper);
L
Liu Jicong 已提交
1846 1847
        taosFreeQitem(pollRspWrapper);
      }
1848
    } else if (pRspWrapper->tmqRspType == TMQ_MSG_TYPE__POLL_META_RSP) {
1849
      // todo handle the wal range and epset for each vgroup
1850
      SMqPollRspWrapper* pollRspWrapper = (SMqPollRspWrapper*)pRspWrapper;
L
Liu Jicong 已提交
1851
      int32_t            consumerEpoch = atomic_load_32(&tmq->epoch);
1852 1853 1854

      tscDebug("consumer:0x%" PRIx64 " process meta rsp", tmq->consumerId);

L
Liu Jicong 已提交
1855
      if (pollRspWrapper->metaRsp.head.epoch == consumerEpoch) {
L
Liu Jicong 已提交
1856
        SMqClientVg* pVg = pollRspWrapper->vgHandle;
H
Haojun Liao 已提交
1857
        pVg->offsetInfo.currentOffset = pollRspWrapper->metaRsp.rspOffset;
L
Liu Jicong 已提交
1858 1859
        atomic_store_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE);
        // build rsp
L
Liu Jicong 已提交
1860
        SMqMetaRspObj* pRsp = tmqBuildMetaRspFromWrapper(pollRspWrapper);
L
Liu Jicong 已提交
1861 1862 1863
        taosFreeQitem(pollRspWrapper);
        return pRsp;
      } else {
H
Haojun Liao 已提交
1864
        tscDebug("consumer:0x%" PRIx64 " vgId:%d msg discard since epoch mismatch: msg epoch %d, consumer epoch %d",
1865
                 tmq->consumerId, pollRspWrapper->vgId, pollRspWrapper->metaRsp.head.epoch, consumerEpoch);
1866
        pRspWrapper = tmqFreeRspWrapper(pRspWrapper);
L
Liu Jicong 已提交
1867
        taosFreeQitem(pollRspWrapper);
X
Xiaoyu Wang 已提交
1868
      }
1869 1870
    } else if (pRspWrapper->tmqRspType == TMQ_MSG_TYPE__TAOSX_RSP) {
      SMqPollRspWrapper* pollRspWrapper = (SMqPollRspWrapper*)pRspWrapper;
X
Xiaoyu Wang 已提交
1871
      int32_t            consumerEpoch = atomic_load_32(&tmq->epoch);
H
Haojun Liao 已提交
1872

L
Liu Jicong 已提交
1873 1874
      if (pollRspWrapper->taosxRsp.head.epoch == consumerEpoch) {
        SMqClientVg* pVg = pollRspWrapper->vgHandle;
H
Haojun Liao 已提交
1875
        pVg->offsetInfo.currentOffset = pollRspWrapper->taosxRsp.rspOffset;
L
Liu Jicong 已提交
1876
        atomic_store_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE);
H
Haojun Liao 已提交
1877

L
Liu Jicong 已提交
1878
        if (pollRspWrapper->taosxRsp.blockNum == 0) {
H
Haojun Liao 已提交
1879 1880
          tscDebug("consumer:0x%" PRIx64 " taosx empty block received, vgId:%d, vg total:%" PRId64 " reqId:0x%" PRIx64,
                   tmq->consumerId, pVg->vgId, pVg->numOfRows, pollRspWrapper->reqId);
H
Haojun Liao 已提交
1881
          pVg->emptyBlockReceiveTs = taosGetTimestampMs();
1882
          pRspWrapper = tmqFreeRspWrapper(pRspWrapper);
H
Haojun Liao 已提交
1883
          taosFreeQitem(pollRspWrapper);
L
Liu Jicong 已提交
1884
          continue;
H
Haojun Liao 已提交
1885
        } else {
X
Xiaoyu Wang 已提交
1886
          pVg->emptyBlockReceiveTs = 0;  // reset the ts
L
Liu Jicong 已提交
1887
        }
wmmhello's avatar
wmmhello 已提交
1888

L
Liu Jicong 已提交
1889
        // build rsp
X
Xiaoyu Wang 已提交
1890
        void*   pRsp = NULL;
1891
        int64_t numOfRows = 0;
L
Liu Jicong 已提交
1892
        if (pollRspWrapper->taosxRsp.createTableNum == 0) {
1893
          pRsp = tmqBuildRspFromWrapper(pollRspWrapper, pVg, &numOfRows);
L
Liu Jicong 已提交
1894
        } else {
wmmhello's avatar
wmmhello 已提交
1895 1896
          pRsp = tmqBuildTaosxRspFromWrapper(pollRspWrapper);
        }
H
Haojun Liao 已提交
1897

1898 1899
        tmq->totalRows += numOfRows;

H
Haojun Liao 已提交
1900
        char buf[80];
H
Haojun Liao 已提交
1901
        tFormatOffset(buf, 80, &pVg->offsetInfo.currentOffset);
H
Haojun Liao 已提交
1902
        tscDebug("consumer:0x%" PRIx64 " process taosx poll rsp, vgId:%d, offset:%s, blocks:%d, rows:%" PRId64
X
Xiaoyu Wang 已提交
1903
                 ", vg total:%" PRId64 " total:%" PRId64 " reqId:0x%" PRIx64,
H
Haojun Liao 已提交
1904
                 tmq->consumerId, pVg->vgId, buf, pollRspWrapper->dataRsp.blockNum, numOfRows, pVg->numOfRows,
H
Haojun Liao 已提交
1905
                 tmq->totalRows, pollRspWrapper->reqId);
H
Haojun Liao 已提交
1906 1907

        taosFreeQitem(pollRspWrapper);
L
Liu Jicong 已提交
1908
        return pRsp;
H
Haojun Liao 已提交
1909

L
Liu Jicong 已提交
1910
      } else {
H
Haojun Liao 已提交
1911
        tscDebug("consumer:0x%" PRIx64 " vgId:%d msg discard since epoch mismatch: msg epoch %d, consumer epoch %d",
1912
                 tmq->consumerId, pollRspWrapper->vgId, pollRspWrapper->taosxRsp.head.epoch, consumerEpoch);
1913
        pRspWrapper = tmqFreeRspWrapper(pRspWrapper);
L
Liu Jicong 已提交
1914 1915
        taosFreeQitem(pollRspWrapper);
      }
X
Xiaoyu Wang 已提交
1916
    } else {
H
Haojun Liao 已提交
1917 1918
      tscDebug("consumer:0x%" PRIx64 " not data msg received", tmq->consumerId);

X
Xiaoyu Wang 已提交
1919
      bool reset = false;
1920 1921
      tmqHandleNoPollRsp(tmq, pRspWrapper, &reset);
      taosFreeQitem(pRspWrapper);
X
Xiaoyu Wang 已提交
1922
      if (pollIfReset && reset) {
1923
        tscDebug("consumer:0x%" PRIx64 ", reset and repoll", tmq->consumerId);
1924
        tmqPollImpl(tmq, timeout);
X
Xiaoyu Wang 已提交
1925 1926 1927 1928 1929
      }
    }
  }
}

1930
TAOS_RES* tmq_consumer_poll(tmq_t* tmq, int64_t timeout) {
L
Liu Jicong 已提交
1931 1932
  void*   rspObj;
  int64_t startTime = taosGetTimestampMs();
L
Liu Jicong 已提交
1933

X
Xiaoyu Wang 已提交
1934 1935
  tscDebug("consumer:0x%" PRIx64 " start to poll at %" PRId64 ", timeout:%" PRId64, tmq->consumerId, startTime,
           timeout);
L
Liu Jicong 已提交
1936

1937
  // in no topic status, delayed task also need to be processed
L
Liu Jicong 已提交
1938
  if (atomic_load_8(&tmq->status) == TMQ_CONSUMER_STATUS__INIT) {
1939
    tscDebug("consumer:0x%" PRIx64 " poll return since consumer is init", tmq->consumerId);
1940
    taosMsleep(500);  //     sleep for a while
1941 1942 1943
    return NULL;
  }

wmmhello's avatar
wmmhello 已提交
1944
  while (atomic_load_8(&tmq->status) == TMQ_CONSUMER_STATUS__RECOVER) {
L
Liu Jicong 已提交
1945
    int32_t retryCnt = 0;
1946
    while (TSDB_CODE_MND_CONSUMER_NOT_READY == doAskEp(tmq)) {
H
Haojun Liao 已提交
1947
      if (retryCnt++ > 40) {
L
Liu Jicong 已提交
1948 1949
        return NULL;
      }
1950

H
Haojun Liao 已提交
1951
      tscDebug("consumer:0x%" PRIx64 " not ready, retry:%d/40 in 500ms", tmq->consumerId, retryCnt);
L
Liu Jicong 已提交
1952 1953 1954 1955
      taosMsleep(500);
    }
  }

X
Xiaoyu Wang 已提交
1956
  while (1) {
L
Liu Jicong 已提交
1957
    tmqHandleAllDelayedTask(tmq);
1958

L
Liu Jicong 已提交
1959
    if (tmqPollImpl(tmq, timeout) < 0) {
1960
      tscDebug("consumer:0x%" PRIx64 " return due to poll error", tmq->consumerId);
L
Liu Jicong 已提交
1961
    }
L
Liu Jicong 已提交
1962

1963
    rspObj = tmqHandleAllRsp(tmq, timeout, false);
L
Liu Jicong 已提交
1964
    if (rspObj) {
1965
      tscDebug("consumer:0x%" PRIx64 " return rsp %p", tmq->consumerId, rspObj);
L
Liu Jicong 已提交
1966
      return (TAOS_RES*)rspObj;
L
Liu Jicong 已提交
1967
    } else if (terrno == TSDB_CODE_TQ_NO_COMMITTED_OFFSET) {
1968
      tscDebug("consumer:0x%" PRIx64 " return null since no committed offset", tmq->consumerId);
L
Liu Jicong 已提交
1969
      return NULL;
X
Xiaoyu Wang 已提交
1970
    }
1971

1972
    if (timeout >= 0) {
L
Liu Jicong 已提交
1973
      int64_t currentTime = taosGetTimestampMs();
1974 1975 1976
      int64_t elapsedTime = currentTime - startTime;
      if (elapsedTime > timeout) {
        tscDebug("consumer:0x%" PRIx64 " (epoch %d) timeout, no rsp, start time %" PRId64 ", current time %" PRId64,
L
Liu Jicong 已提交
1977
                 tmq->consumerId, tmq->epoch, startTime, currentTime);
X
Xiaoyu Wang 已提交
1978 1979
        return NULL;
      }
1980
      tsem_timewait(&tmq->rspSem, (timeout - elapsedTime));
L
Liu Jicong 已提交
1981 1982
    } else {
      // use tsem_timewait instead of tsem_wait to avoid unexpected stuck
L
Liu Jicong 已提交
1983
      tsem_timewait(&tmq->rspSem, 1000);
X
Xiaoyu Wang 已提交
1984 1985 1986 1987
    }
  }
}

1988 1989 1990 1991 1992 1993 1994 1995 1996 1997 1998 1999 2000 2001
static void displayConsumeStatistics(const tmq_t* pTmq) {
  int32_t numOfTopics = taosArrayGetSize(pTmq->clientTopics);
  tscDebug("consumer:0x%" PRIx64 " closing poll:%" PRId64 " rows:%" PRId64 " topics:%d, final epoch:%d",
           pTmq->consumerId, pTmq->pollCnt, pTmq->totalRows, numOfTopics, pTmq->epoch);

  tscDebug("consumer:0x%" PRIx64 " rows dist begin: ", pTmq->consumerId);
  for (int32_t i = 0; i < numOfTopics; ++i) {
    SMqClientTopic* pTopics = taosArrayGet(pTmq->clientTopics, i);

    tscDebug("consumer:0x%" PRIx64 " topic:%d", pTmq->consumerId, i);
    int32_t numOfVgs = taosArrayGetSize(pTopics->vgs);
    for (int32_t j = 0; j < numOfVgs; ++j) {
      SMqClientVg* pVg = taosArrayGet(pTopics->vgs, j);
      tscDebug("topic:%s, %d. vgId:%d rows:%" PRId64, pTopics->topicName, j, pVg->vgId, pVg->numOfRows);
2002
    }
2003
  }
2004

2005 2006
  tscDebug("consumer:0x%" PRIx64 " rows dist end", pTmq->consumerId);
}
2007

2008
int32_t tmq_consumer_close(tmq_t* tmq) {
X
Xiaoyu Wang 已提交
2009
  tscDebug("consumer:0x%" PRIx64 " start to close consumer, status:%d", tmq->consumerId, tmq->status);
2010
  displayConsumeStatistics(tmq);
2011

2012 2013 2014 2015 2016 2017
  if (tmq->status == TMQ_CONSUMER_STATUS__READY) {
    // if auto commit is set, commit before close consumer. Otherwise, do nothing.
    if (tmq->autoCommit) {
      int32_t rsp = tmq_commit_sync(tmq, NULL);
      if (rsp != 0) {
        return rsp;
2018 2019 2020
      }
    }

L
Liu Jicong 已提交
2021
    int32_t     retryCnt = 0;
2022
    tmq_list_t* lst = tmq_list_new();
L
Liu Jicong 已提交
2023
    while (1) {
2024
      int32_t rsp = tmq_subscribe(tmq, lst);
L
Liu Jicong 已提交
2025 2026 2027 2028 2029 2030 2031 2032
      if (rsp != TSDB_CODE_MND_CONSUMER_NOT_READY || retryCnt > 5) {
        break;
      } else {
        retryCnt++;
        taosMsleep(500);
      }
    }

2033
    tmq_list_destroy(lst);
2034
  } else {
X
Xiaoyu Wang 已提交
2035
    tscWarn("consumer:0x%" PRIx64 " not in ready state, close it directly", tmq->consumerId);
L
Liu Jicong 已提交
2036
  }
H
Haojun Liao 已提交
2037

2038
  taosRemoveRef(tmqMgmt.rsetId, tmq->refId);
L
Liu Jicong 已提交
2039
  return 0;
2040
}
L
Liu Jicong 已提交
2041

L
Liu Jicong 已提交
2042 2043
const char* tmq_err2str(int32_t err) {
  if (err == 0) {
L
Liu Jicong 已提交
2044
    return "success";
L
Liu Jicong 已提交
2045
  } else if (err == -1) {
L
Liu Jicong 已提交
2046 2047 2048
    return "fail";
  } else {
    return tstrerror(err);
L
Liu Jicong 已提交
2049 2050
  }
}
L
Liu Jicong 已提交
2051

L
Liu Jicong 已提交
2052 2053 2054 2055 2056
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;
2057 2058
  } else if (TD_RES_TMQ_METADATA(res)) {
    return TMQ_RES_METADATA;
L
Liu Jicong 已提交
2059 2060 2061 2062 2063
  } else {
    return TMQ_RES_INVALID;
  }
}

L
Liu Jicong 已提交
2064
const char* tmq_get_topic_name(TAOS_RES* res) {
L
Liu Jicong 已提交
2065 2066
  if (TD_RES_TMQ(res)) {
    SMqRspObj* pRspObj = (SMqRspObj*)res;
L
Liu Jicong 已提交
2067
    return strchr(pRspObj->topic, '.') + 1;
L
Liu Jicong 已提交
2068 2069 2070
  } else if (TD_RES_TMQ_META(res)) {
    SMqMetaRspObj* pMetaRspObj = (SMqMetaRspObj*)res;
    return strchr(pMetaRspObj->topic, '.') + 1;
2071 2072 2073
  } else if (TD_RES_TMQ_METADATA(res)) {
    SMqTaosxRspObj* pRspObj = (SMqTaosxRspObj*)res;
    return strchr(pRspObj->topic, '.') + 1;
L
Liu Jicong 已提交
2074 2075 2076 2077 2078
  } else {
    return NULL;
  }
}

L
Liu Jicong 已提交
2079 2080 2081 2082
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 已提交
2083 2084 2085
  } else if (TD_RES_TMQ_META(res)) {
    SMqMetaRspObj* pMetaRspObj = (SMqMetaRspObj*)res;
    return strchr(pMetaRspObj->db, '.') + 1;
2086 2087 2088
  } else if (TD_RES_TMQ_METADATA(res)) {
    SMqTaosxRspObj* pRspObj = (SMqTaosxRspObj*)res;
    return strchr(pRspObj->db, '.') + 1;
L
Liu Jicong 已提交
2089 2090 2091 2092 2093
  } else {
    return NULL;
  }
}

L
Liu Jicong 已提交
2094 2095 2096 2097
int32_t tmq_get_vgroup_id(TAOS_RES* res) {
  if (TD_RES_TMQ(res)) {
    SMqRspObj* pRspObj = (SMqRspObj*)res;
    return pRspObj->vgId;
L
Liu Jicong 已提交
2098 2099 2100
  } else if (TD_RES_TMQ_META(res)) {
    SMqMetaRspObj* pMetaRspObj = (SMqMetaRspObj*)res;
    return pMetaRspObj->vgId;
2101
  } else if (TD_RES_TMQ_METADATA(res)) {
2102 2103
    SMqTaosxRspObj* pRspObj = (SMqTaosxRspObj*)res;
    return pRspObj->vgId;
L
Liu Jicong 已提交
2104 2105 2106 2107
  } else {
    return -1;
  }
}
L
Liu Jicong 已提交
2108 2109 2110 2111 2112 2113 2114 2115

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;
    }
2116
    return (const char*)taosArrayGetP(pRspObj->rsp.blockTbName, pRspObj->resIter);
2117 2118
  } else if (TD_RES_TMQ_METADATA(res)) {
    SMqTaosxRspObj* pRspObj = (SMqTaosxRspObj*)res;
L
Liu Jicong 已提交
2119 2120 2121
    if (!pRspObj->rsp.withTbName || pRspObj->rsp.blockTbName == NULL || pRspObj->resIter < 0 ||
        pRspObj->resIter >= pRspObj->rsp.blockNum) {
      return NULL;
2122
    }
L
Liu Jicong 已提交
2123 2124
    return (const char*)taosArrayGetP(pRspObj->rsp.blockTbName, pRspObj->resIter);
  }
L
Liu Jicong 已提交
2125 2126
  return NULL;
}
2127

2128 2129 2130 2131
void tmq_commit_async(tmq_t* tmq, const TAOS_RES* pRes, tmq_commit_cb* cb, void* param) {
  if (pRes == NULL) {  // here needs to commit all offsets.
    asyncCommitAllOffsets(tmq, cb, param);
  } else {  // only commit one offset
2132
    asyncCommitOffset(tmq, pRes, TDMT_VND_TMQ_COMMIT_OFFSET, cb, param);
2133
  }
L
Liu Jicong 已提交
2134 2135
}

2136
static void commitCallBackFn(tmq_t *UNUSED_PARAM(tmq), int32_t code, void* param) {
2137 2138 2139
  SSyncCommitInfo* pInfo = (SSyncCommitInfo*) param;
  pInfo->code = code;
  tsem_post(&pInfo->sem);
2140
}
2141

2142 2143 2144 2145 2146 2147 2148 2149 2150
int32_t tmq_commit_sync(tmq_t* tmq, const TAOS_RES* pRes) {
  int32_t code = 0;

  SSyncCommitInfo* pInfo = taosMemoryMalloc(sizeof(SSyncCommitInfo));
  tsem_init(&pInfo->sem, 0, 0);
  pInfo->code = 0;

  if (pRes == NULL) {
    asyncCommitAllOffsets(tmq, commitCallBackFn, pInfo);
H
Haojun Liao 已提交
2151
  } else {
2152
    asyncCommitOffset(tmq, pRes, TDMT_VND_TMQ_COMMIT_OFFSET, commitCallBackFn, pInfo);
2153 2154
  }

2155 2156
  tsem_wait(&pInfo->sem);
  code = pInfo->code;
H
Haojun Liao 已提交
2157 2158

  tsem_destroy(&pInfo->sem);
2159 2160
  taosMemoryFree(pInfo);

X
Xiaoyu Wang 已提交
2161
  tscDebug("consumer:0x%" PRIx64 " sync commit done, code:%s", tmq->consumerId, tstrerror(code));
2162 2163 2164 2165 2166 2167 2168 2169 2170 2171 2172 2173
  return code;
}

void updateEpCallbackFn(tmq_t* pTmq, int32_t code, SDataBuf* pDataBuf, void* param) {
  SAskEpInfo* pInfo = param;
  pInfo->code = code;

  if (code == TSDB_CODE_SUCCESS) {
    SMqRspHead* head = pDataBuf->pData;

    SMqAskEpRsp rsp;
    tDecodeSMqAskEpRsp(POINTER_SHIFT(pDataBuf->pData, sizeof(SMqRspHead)), &rsp);
2174
    doUpdateLocalEp(pTmq, head->epoch, &rsp);
2175 2176 2177
    tDeleteSMqAskEpRsp(&rsp);
  }

H
Haojun Liao 已提交
2178
  tsem_post(&pInfo->sem);
2179 2180 2181 2182 2183 2184 2185 2186 2187 2188 2189 2190 2191 2192 2193 2194 2195 2196 2197 2198 2199 2200 2201 2202 2203 2204
}

void addToQueueCallbackFn(tmq_t* pTmq, int32_t code, SDataBuf* pDataBuf, void* param) {
  if (code != TSDB_CODE_SUCCESS) {
    terrno = code;
    return;
  }

  SMqAskEpRspWrapper* pWrapper = taosAllocateQitem(sizeof(SMqAskEpRspWrapper), DEF_QITEM, 0);
  if (pWrapper == NULL) {
    terrno = TSDB_CODE_OUT_OF_MEMORY;
    return;
  }

  SMqRspHead* head = pDataBuf->pData;

  pWrapper->tmqRspType = TMQ_MSG_TYPE__EP_RSP;
  pWrapper->epoch = head->epoch;
  memcpy(&pWrapper->msg, pDataBuf->pData, sizeof(SMqRspHead));
  tDecodeSMqAskEpRsp(POINTER_SHIFT(pDataBuf->pData, sizeof(SMqRspHead)), &pWrapper->msg);

  taosWriteQitem(pTmq->mqueue, pWrapper);
}

int32_t doAskEp(tmq_t* pTmq) {
  SAskEpInfo* pInfo = taosMemoryMalloc(sizeof(SAskEpInfo));
H
Haojun Liao 已提交
2205
  tsem_init(&pInfo->sem, 0, 0);
2206 2207

  asyncAskEp(pTmq, updateEpCallbackFn, pInfo);
H
Haojun Liao 已提交
2208
  tsem_wait(&pInfo->sem);
2209 2210

  int32_t code = pInfo->code;
H
Haojun Liao 已提交
2211
  tsem_destroy(&pInfo->sem);
2212 2213 2214 2215 2216
  taosMemoryFree(pInfo);
  return code;
}

void asyncAskEp(tmq_t* pTmq, __tmq_askep_fn_t askEpFn, void* param) {
2217
  SMqAskEpReq req = {0};
2218 2219 2220
  req.consumerId = pTmq->consumerId;
  req.epoch = pTmq->epoch;
  strcpy(req.cgroup, pTmq->groupId);
2221 2222 2223

  int32_t tlen = tSerializeSMqAskEpReq(NULL, 0, &req);
  if (tlen < 0) {
2224 2225 2226
    tscError("consumer:0x%" PRIx64 ", tSerializeSMqAskEpReq failed", pTmq->consumerId);
    askEpFn(pTmq, TSDB_CODE_INVALID_PARA, NULL, param);
    return;
2227 2228 2229 2230
  }

  void* pReq = taosMemoryCalloc(1, tlen);
  if (pReq == NULL) {
2231 2232 2233
    tscError("consumer:0x%" PRIx64 ", failed to malloc askEpReq msg, size:%d", pTmq->consumerId, tlen);
    askEpFn(pTmq, TSDB_CODE_OUT_OF_MEMORY, NULL, param);
    return;
2234 2235 2236
  }

  if (tSerializeSMqAskEpReq(pReq, tlen, &req) < 0) {
2237
    tscError("consumer:0x%" PRIx64 ", tSerializeSMqAskEpReq %d failed", pTmq->consumerId, tlen);
2238
    taosMemoryFree(pReq);
2239 2240 2241

    askEpFn(pTmq, TSDB_CODE_INVALID_PARA, NULL, param);
    return;
2242 2243 2244 2245
  }

  SMqAskEpCbParam* pParam = taosMemoryCalloc(1, sizeof(SMqAskEpCbParam));
  if (pParam == NULL) {
2246
    tscError("consumer:0x%" PRIx64 ", failed to malloc subscribe param", pTmq->consumerId);
2247
    taosMemoryFree(pReq);
2248 2249 2250

    askEpFn(pTmq, TSDB_CODE_OUT_OF_MEMORY, NULL, param);
    return;
2251 2252
  }

2253 2254 2255 2256
  pParam->refId = pTmq->refId;
  pParam->epoch = pTmq->epoch;
  pParam->pUserFn = askEpFn;
  pParam->pParam = param;
2257 2258 2259 2260 2261

  SMsgSendInfo* sendInfo = taosMemoryCalloc(1, sizeof(SMsgSendInfo));
  if (sendInfo == NULL) {
    taosMemoryFree(pParam);
    taosMemoryFree(pReq);
2262 2263
    askEpFn(pTmq, TSDB_CODE_OUT_OF_MEMORY, NULL, param);
    return;
2264 2265
  }

X
Xiaoyu Wang 已提交
2266
  sendInfo->msgInfo = (SDataBuf){.pData = pReq, .len = tlen, .handle = NULL};
2267 2268 2269 2270

  sendInfo->requestId = generateRequestId();
  sendInfo->requestObjRefId = 0;
  sendInfo->param = pParam;
2271
  sendInfo->fp = askEpCallbackFn;
2272 2273
  sendInfo->msgType = TDMT_MND_TMQ_ASK_EP;

2274 2275
  SEpSet epSet = getEpSet_s(&pTmq->pTscObj->pAppInfo->mgmtEp);
  tscDebug("consumer:0x%" PRIx64 " ask ep from mnode, reqId:0x%" PRIx64, pTmq->consumerId, sendInfo->requestId);
2276 2277

  int64_t transporterId = 0;
2278
  asyncSendMsgToServer(pTmq->pTscObj->pAppInfo->pTransporter, &epSet, &transporterId, sendInfo);
2279 2280 2281 2282 2283 2284 2285
}

int32_t makeTopicVgroupKey(char* dst, const char* topicName, int32_t vg) {
  return sprintf(dst, "%s:%d", topicName, vg);
}

int32_t tmqCommitDone(SMqCommitCbParamSet* pParamSet) {
2286 2287 2288
  int64_t refId = pParamSet->refId;

  tmq_t* tmq = taosAcquireRef(tmqMgmt.rsetId, refId);
2289 2290 2291 2292 2293 2294 2295
  if (tmq == NULL) {
    taosMemoryFree(pParamSet);
    terrno = TSDB_CODE_TMQ_CONSUMER_CLOSED;
    return -1;
  }

  // if no more waiting rsp
2296
  pParamSet->callbackFn(tmq, pParamSet->code, pParamSet->userParam);
2297
  taosMemoryFree(pParamSet);
2298 2299

  taosReleaseRef(tmqMgmt.rsetId, refId);
2300
  return 0;
2301 2302
}

2303
void commitRspCountDown(SMqCommitCbParamSet* pParamSet, int64_t consumerId, const char* pTopic, int32_t vgId) {
2304 2305
  int32_t waitingRspNum = atomic_sub_fetch_32(&pParamSet->waitingRspNum, 1);
  if (waitingRspNum == 0) {
H
Haojun Liao 已提交
2306 2307
    tscDebug("consumer:0x%" PRIx64 " topic:%s vgId:%d all commit-rsp received, commit completed", consumerId, pTopic,
             vgId);
2308
    tmqCommitDone(pParamSet);
H
Haojun Liao 已提交
2309 2310 2311
  } else {
    tscDebug("consumer:0x%" PRIx64 " topic:%s vgId:%d commit-rsp received, remain:%d", consumerId, pTopic, vgId,
             waitingRspNum);
2312 2313
  }
}
2314 2315 2316 2317 2318 2319 2320 2321 2322 2323 2324 2325 2326 2327 2328 2329 2330 2331 2332 2333 2334 2335

SReqResultInfo* tmqGetNextResInfo(TAOS_RES* res, bool convertUcs4) {
  SMqRspObj* pRspObj = (SMqRspObj*)res;
  pRspObj->resIter++;

  if (pRspObj->resIter < pRspObj->rsp.blockNum) {
    SRetrieveTableRsp* pRetrieve = (SRetrieveTableRsp*)taosArrayGetP(pRspObj->rsp.blockData, pRspObj->resIter);
    if (pRspObj->rsp.withSchema) {
      SSchemaWrapper* pSW = (SSchemaWrapper*)taosArrayGetP(pRspObj->rsp.blockSchema, pRspObj->resIter);
      setResSchemaInfo(&pRspObj->resInfo, pSW->pSchema, pSW->nCols);
      taosMemoryFreeClear(pRspObj->resInfo.row);
      taosMemoryFreeClear(pRspObj->resInfo.pCol);
      taosMemoryFreeClear(pRspObj->resInfo.length);
      taosMemoryFreeClear(pRspObj->resInfo.convertBuf);
      taosMemoryFreeClear(pRspObj->resInfo.convertJson);
    }

    setQueryResultFromRsp(&pRspObj->resInfo, pRetrieve, convertUcs4, false);
    return &pRspObj->resInfo;
  }

  return NULL;
H
Haojun Liao 已提交
2336 2337
}

2338 2339 2340 2341 2342 2343 2344 2345 2346 2347 2348 2349 2350 2351 2352 2353 2354 2355 2356 2357 2358 2359 2360 2361 2362 2363 2364 2365 2366 2367 2368 2369 2370 2371 2372 2373 2374 2375 2376 2377 2378 2379 2380 2381
static int32_t tmqGetWalInfoCb(void* param, SDataBuf* pMsg, int32_t code) {
  SMqVgWalInfoParam* pParam = param;
  SMqVgCommon* pCommon = pParam->pCommon;

  int32_t total = atomic_add_fetch_32(&pCommon->numOfRsp, 1);
  if (code != TSDB_CODE_SUCCESS) {
    tscError("consumer:0x%" PRIx64 " failed to get the wal info from vgId:%d for topic:%s", pCommon->consumerId,
             pParam->vgId, pCommon->pTopicName);
    pCommon->code = code;
  } else {
    SMqDataRsp rsp;
    SDecoder decoder;
    tDecoderInit(&decoder, POINTER_SHIFT(pMsg->pData, sizeof(SMqRspHead)), pMsg->len - sizeof(SMqRspHead));
    tDecodeMqDataRsp(&decoder, &rsp);
    tDecoderClear(&decoder);

    SMqRspHead* pHead = pMsg->pData;

    tmq_topic_assignment assignment = {.begin = pHead->walsver,
                                       .end = pHead->walever,
                                       .currentOffset = rsp.rspOffset.version,
                                       .vgroupHandle = pParam->vgId};

    taosThreadMutexLock(&pCommon->mutex);
    taosArrayPush(pCommon->pList, &assignment);
    taosThreadMutexUnlock(&pCommon->mutex);
  }

  if (total == pParam->totalReq) {
    tsem_post(&pCommon->rsp);
  }

  taosMemoryFree(pParam);
  return 0;
}

static void destroyCommonInfo(SMqVgCommon* pCommon) {
  taosArrayDestroy(pCommon->pList);
  tsem_destroy(&pCommon->rsp);
  taosThreadMutexDestroy(&pCommon->mutex);
  taosMemoryFree(pCommon->pTopicName);
  taosMemoryFree(pCommon);
}

H
Haojun Liao 已提交
2382
int32_t tmq_get_topic_assignment(tmq_t* tmq, const char* pTopicName, tmq_topic_assignment** assignment,
H
Haojun Liao 已提交
2383
                                 int32_t* numOfAssignment) {
H
Haojun Liao 已提交
2384 2385 2386
  *numOfAssignment = 0;
  *assignment = NULL;

2387
  int32_t accId = tmq->pTscObj->acctId;
2388
  char    tname[128] = {0};
2389 2390 2391
  sprintf(tname, "%d.%s", accId, pTopicName);

  SMqClientTopic* pTopic = getTopicByName(tmq, tname);
H
Haojun Liao 已提交
2392 2393 2394 2395 2396 2397
  if (pTopic == NULL) {
    return TSDB_CODE_INVALID_PARA;
  }

  // in case of snapshot is opened, no valid offset will return
  *numOfAssignment = taosArrayGetSize(pTopic->vgs);
2398 2399 2400 2401 2402 2403 2404 2405

  *assignment = taosMemoryCalloc(*numOfAssignment, sizeof(tmq_topic_assignment));
  if (*assignment == NULL) {
    tscError("consumer:0x%" PRIx64 " failed to malloc buffer, size:%" PRIzu, tmq->consumerId,
             (*numOfAssignment) * sizeof(tmq_topic_assignment));
    return TSDB_CODE_OUT_OF_MEMORY;
  }

2406 2407
  bool needFetch = false;

H
Haojun Liao 已提交
2408 2409
  for (int32_t j = 0; j < (*numOfAssignment); ++j) {
    SMqClientVg* pClientVg = taosArrayGet(pTopic->vgs, j);
2410 2411 2412 2413
    if (!pClientVg->receiveInfo) {
      needFetch = true;
      break;
    }
H
Haojun Liao 已提交
2414 2415 2416 2417 2418 2419 2420 2421 2422 2423 2424 2425 2426

    tmq_topic_assignment* pAssignment = &(*assignment)[j];
    if (pClientVg->offsetInfo.currentOffset.type == TMQ_OFFSET__LOG) {
      pAssignment->currentOffset = pClientVg->offsetInfo.currentOffset.version;
    } else {
      pAssignment->currentOffset = 0;
    }

    pAssignment->begin = pClientVg->offsetInfo.walVerBegin;
    pAssignment->end = pClientVg->offsetInfo.walVerEnd;
    pAssignment->vgroupHandle = pClientVg->vgId;
  }

2427 2428 2429 2430 2431 2432 2433 2434 2435 2436 2437 2438 2439 2440 2441 2442 2443 2444 2445 2446 2447 2448 2449 2450 2451 2452 2453 2454 2455 2456 2457 2458 2459 2460 2461 2462 2463 2464 2465 2466 2467 2468 2469 2470 2471 2472 2473 2474 2475 2476 2477 2478 2479 2480 2481 2482 2483 2484 2485 2486 2487 2488 2489 2490 2491 2492 2493 2494 2495 2496 2497 2498 2499 2500 2501 2502 2503 2504 2505 2506 2507 2508 2509 2510 2511 2512 2513 2514 2515 2516 2517 2518 2519 2520 2521 2522
  if (needFetch) {
    SMqVgCommon* pCommon = taosMemoryCalloc(1, sizeof(SMqVgCommon));
    if (pCommon == NULL) {
      terrno = TSDB_CODE_OUT_OF_MEMORY;
      return terrno;
    }

    pCommon->pList= taosArrayInit(4, sizeof(tmq_topic_assignment));
    tsem_init(&pCommon->rsp, 0, 0);
    taosThreadMutexInit(&pCommon->mutex, 0);
    pCommon->pTopicName = taosStrdup(pTopic->topicName);
    pCommon->consumerId = tmq->consumerId;

    terrno = TSDB_CODE_OUT_OF_MEMORY;
    for (int32_t i = 0; i < (*numOfAssignment); ++i) {
      SMqClientVg* pClientVg = taosArrayGet(pTopic->vgs, i);

      SMqVgWalInfoParam* pParam = taosMemoryMalloc(sizeof(SMqVgWalInfoParam));
      if (pParam == NULL) {
        destroyCommonInfo(pCommon);
        return terrno;
      }

      pParam->epoch = tmq->epoch;
      pParam->vgId = pClientVg->vgId;
      pParam->totalReq = *numOfAssignment;
      pParam->pCommon = pCommon;

      SMqPollReq req = {0};
      tmqBuildConsumeReqImpl(&req, tmq, 10, pTopic, pClientVg);

      int32_t msgSize = tSerializeSMqPollReq(NULL, 0, &req);
      if (msgSize < 0) {
        taosMemoryFree(pParam);
        destroyCommonInfo(pCommon);
        return terrno;
      }

      char* msg = taosMemoryCalloc(1, msgSize);
      if (NULL == msg) {
        taosMemoryFree(pParam);
        destroyCommonInfo(pCommon);
        return terrno;
      }

      if (tSerializeSMqPollReq(msg, msgSize, &req) < 0) {
        taosMemoryFree(msg);
        taosMemoryFree(pParam);
        destroyCommonInfo(pCommon);
        return terrno;
      }

      SMsgSendInfo* sendInfo = taosMemoryCalloc(1, sizeof(SMsgSendInfo));
      if (sendInfo == NULL) {
        taosMemoryFree(pParam);
        taosMemoryFree(msg);
        destroyCommonInfo(pCommon);
        return terrno;
      }

      sendInfo->msgInfo = (SDataBuf){.pData = msg, .len = msgSize, .handle = NULL};
      sendInfo->requestId = req.reqId;
      sendInfo->requestObjRefId = 0;
      sendInfo->param = pParam;
      sendInfo->fp = tmqGetWalInfoCb;
      sendInfo->msgType = TDMT_VND_TMQ_VG_WALINFO;

      int64_t transporterId = 0;
      char    offsetFormatBuf[80];
      tFormatOffset(offsetFormatBuf, tListLen(offsetFormatBuf), &pClientVg->offsetInfo.currentOffset);

      tscDebug("consumer:0x%" PRIx64 " %s retrieve wal info vgId:%d, epoch %d, req:%s, reqId:0x%" PRIx64,
               tmq->consumerId, pTopic->topicName, pClientVg->vgId, tmq->epoch, offsetFormatBuf, req.reqId);
      asyncSendMsgToServer(tmq->pTscObj->pAppInfo->pTransporter, &pClientVg->epSet, &transporterId, sendInfo);
    }

    tsem_wait(&pCommon->rsp);
    int32_t code = pCommon->code;

    terrno = code;
    if (code != TSDB_CODE_SUCCESS) {
      taosMemoryFree(*assignment);
      *numOfAssignment = 0;
    } else {
      int32_t num = taosArrayGetSize(pCommon->pList);
      for(int32_t i = 0; i < num; ++i) {
        (*assignment)[i] = *(tmq_topic_assignment*)taosArrayGet(pCommon->pList, i);
      }
      *numOfAssignment = num;
    }

    destroyCommonInfo(pCommon);
    return code;
  } else {
    return TSDB_CODE_SUCCESS;
  }
H
Haojun Liao 已提交
2523 2524 2525 2526
}

int32_t tmq_offset_seek(tmq_t* tmq, const char* pTopicName, int32_t vgroupHandle, int64_t offset) {
  if (tmq == NULL) {
H
Haojun Liao 已提交
2527
    tscError("invalid tmq handle, null");
H
Haojun Liao 已提交
2528 2529 2530
    return TSDB_CODE_INVALID_PARA;
  }

2531 2532 2533 2534 2535
  int32_t accId = tmq->pTscObj->acctId;
  char tname[128] = {0};
  sprintf(tname, "%d.%s", accId, pTopicName);

  SMqClientTopic* pTopic = getTopicByName(tmq, tname);
H
Haojun Liao 已提交
2536
  if (pTopic == NULL) {
2537
    tscError("consumer:0x%" PRIx64 " invalid topic name:%s", tmq->consumerId, pTopicName);
H
Haojun Liao 已提交
2538 2539 2540 2541
    return TSDB_CODE_INVALID_PARA;
  }

  SMqClientVg* pVg = NULL;
H
Haojun Liao 已提交
2542 2543
  int32_t      numOfVgs = taosArrayGetSize(pTopic->vgs);
  for (int32_t i = 0; i < numOfVgs; ++i) {
H
Haojun Liao 已提交
2544 2545 2546 2547 2548 2549 2550 2551
    SMqClientVg* pClientVg = taosArrayGet(pTopic->vgs, i);
    if (pClientVg->vgId == vgroupHandle) {
      pVg = pClientVg;
      break;
    }
  }

  if (pVg == NULL) {
2552
    tscError("consumer:0x%" PRIx64 " invalid vgroup id:%d", tmq->consumerId, vgroupHandle);
H
Haojun Liao 已提交
2553 2554 2555 2556 2557 2558 2559
    return TSDB_CODE_INVALID_PARA;
  }

  SVgOffsetInfo* pOffsetInfo = &pVg->offsetInfo;

  int32_t type = pOffsetInfo->currentOffset.type;
  if (type != TMQ_OFFSET__LOG) {
2560
    tscError("consumer:0x%" PRIx64 " offset type:%d not wal version, seek not allowed", tmq->consumerId, type);
H
Haojun Liao 已提交
2561 2562 2563
    return TSDB_CODE_INVALID_PARA;
  }

H
Haojun Liao 已提交
2564
  if (offset < pOffsetInfo->walVerBegin || offset > pOffsetInfo->walVerEnd) {
2565
    tscError("consumer:0x%" PRIx64 " invalid seek params, offset:%" PRId64, tmq->consumerId, offset);
H
Haojun Liao 已提交
2566 2567 2568
    return TSDB_CODE_INVALID_PARA;
  }

H
Haojun Liao 已提交
2569 2570 2571
  // update the offset, and then commit to vnode
  if (pOffsetInfo->currentOffset.type == TMQ_OFFSET__LOG) {
    pOffsetInfo->currentOffset.version = offset;
2572
    pOffsetInfo->committedOffset.version = INT64_MIN;
H
Haojun Liao 已提交
2573 2574
  }

H
Haojun Liao 已提交
2575
  SMqRspObj rspObj = {.resType = RES_TYPE__TMQ, .vgId = pVg->vgId};
2576
  tstrncpy(rspObj.topic, tname, tListLen(rspObj.topic));
H
Haojun Liao 已提交
2577

H
Haojun Liao 已提交
2578
  tscDebug("consumer:0x%" PRIx64 " seek to %" PRId64 " on vgId:%d", tmq->consumerId, offset, pVg->vgId);
2579 2580 2581 2582 2583 2584 2585 2586 2587 2588 2589 2590 2591 2592 2593 2594 2595 2596

  SSyncCommitInfo* pInfo = taosMemoryMalloc(sizeof(SSyncCommitInfo));
  if (pInfo == NULL) {
    tscError("consumer:0x%"PRIx64" failed to prepare seek operation", tmq->consumerId);
    return TSDB_CODE_OUT_OF_MEMORY;
  }

  tsem_init(&pInfo->sem, 0, 0);
  pInfo->code = 0;

  asyncCommitOffset(tmq, &rspObj, TDMT_VND_TMQ_SEEK_TO_OFFSET, commitCallBackFn, pInfo);

  tsem_wait(&pInfo->sem);
  int32_t code = pInfo->code;

  tsem_destroy(&pInfo->sem);
  taosMemoryFree(pInfo);

H
Haojun Liao 已提交
2597 2598 2599 2600 2601 2602
  if (code != TSDB_CODE_SUCCESS) {
    tscError("consumer:0x%" PRIx64 " failed to send seek to vgId:%d, code:%s", tmq->consumerId, pVg->vgId,
             tstrerror(code));
  }

  return code;
2603
}