clientTmq.c 85.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
} SMqCommitCbParamSet;

typedef struct {
  SMqCommitCbParamSet* params;
231
  SMqVgOffset*         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);
L
fix  
Liu Jicong 已提交
400
  if (taosArrayPush(container, &topic) == NULL) return -1;
401 402 403
  return 0;
}

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

L
Liu Jicong 已提交
409 410 411 412 413 414 415 416 417 418
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 已提交
419 420
static SMqClientVg* foundClientVg(SArray* pTopicList, const char* pName, int32_t vgId, int32_t* index,
                                  int32_t* numOfVgroups) {
421 422 423 424
  int32_t numOfTopics = taosArrayGetSize(pTopicList);
  *index = -1;
  *numOfVgroups = 0;

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

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

  return NULL;
L
Liu Jicong 已提交
442
}
443

444 445 446
// 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 已提交
447
static int32_t tmqCommitCb(void* param, SDataBuf* pBuf, int32_t code) {
448
  SMqCommitCbParam*    pParam = (SMqCommitCbParam*)param;
449
  SMqCommitCbParamSet* pParamSet = (SMqCommitCbParamSet*)pParam->params;
H
Haojun Liao 已提交
450

X
Xiaoyu Wang 已提交
451 452 453 454 455 456 457 458 459 460 461 462 463 464
  //  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 已提交
465
  //                 pParam->pTmq->consumerId, pParam->topicName, pVg->vgId, pVg->offsetInfo.committedOffset.version,
X
Xiaoyu Wang 已提交
466 467 468 469 470 471 472 473 474 475 476 477 478 479 480
  //                 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
481

L
Liu Jicong 已提交
482
  taosMemoryFree(pParam->pOffset);
L
Liu Jicong 已提交
483
  taosMemoryFree(pBuf->pData);
dengyihao's avatar
dengyihao 已提交
484
  taosMemoryFree(pBuf->pEpSet);
L
Liu Jicong 已提交
485

486
  commitRspCountDown(pParamSet, pParam->pTmq->consumerId, pParam->topicName, pParam->vgId);
487 488 489
  return 0;
}

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

497 498
  pOffset->consumerId = tmq->consumerId;
  pOffset->offset.val = pVg->offsetInfo.currentOffset;
499

L
Liu Jicong 已提交
500
  int32_t groupLen = strlen(tmq->groupId);
501 502 503
  memcpy(pOffset->offset.subKey, tmq->groupId, groupLen);
  pOffset->offset.subKey[groupLen] = TMQ_SEPARATOR;
  strcpy(pOffset->offset.subKey + groupLen + 1, pTopicName);
L
Liu Jicong 已提交
504

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

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

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

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

  SEncoder encoder;
  tEncoderInit(&encoder, abuf, len);
524
  tEncodeMqVgOffset(&encoder, pOffset);
L
Liu Jicong 已提交
525
  tEncoderClear(&encoder);
L
Liu Jicong 已提交
526 527

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

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

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

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

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

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

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

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

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

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

  return TSDB_CODE_SUCCESS;
L
Liu Jicong 已提交
577 578
}

H
Haojun Liao 已提交
579 580 581 582 583 584 585 586 587 588 589
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 已提交
590
  tscError("consumer:0x%" PRIx64 ", total:%d, failed to find topic:%s", tmq->consumerId, numOfTopics, pTopicName);
H
Haojun Liao 已提交
591 592 593
  return NULL;
}

594
static void asyncCommitOffset(tmq_t* tmq, const TAOS_RES* pRes, int32_t type, tmq_commit_cb* pCommitFp, void* userParam) {
595 596 597 598 599 600 601 602 603 604 605 606
  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 已提交
607
    vgId = pRspObj->vgId;
608 609 610
  } else if (TD_RES_TMQ_META(pRes)) {
    SMqMetaRspObj* pMetaRspObj = (SMqMetaRspObj*)pRes;
    pTopicName = pMetaRspObj->topic;
L
Liu Jicong 已提交
611
    vgId = pMetaRspObj->vgId;
612 613 614
  } else if (TD_RES_TMQ_METADATA(pRes)) {
    SMqTaosxRspObj* pRspObj = (SMqTaosxRspObj*)pRes;
    pTopicName = pRspObj->topic;
615
    vgId = pRspObj->vgId;
L
Liu Jicong 已提交
616
  } else {
617 618
    pCommitFp(tmq, TSDB_CODE_TMQ_INVALID_MSG, userParam);
    return;
L
Liu Jicong 已提交
619 620 621 622
  }

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

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

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

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

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

645 646 647 648 649 650
  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 已提交
651
    }
L
Liu Jicong 已提交
652
  }
L
Liu Jicong 已提交
653

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

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

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

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

  pParamSet->refId = tmq->refId;
  pParamSet->epoch = tmq->epoch;
686
  pParamSet->callbackFn = pCommitFp;
687 688
  pParamSet->userParam = userParam;

689 690 691
  // init as 1 to prevent concurrency issue
  pParamSet->waitingRspNum = 1;

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

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

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

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

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

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

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

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

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

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

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

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

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

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

void tmqSendHbReq(void* param, void* tmrId) {
782
  int64_t refId = *(int64_t*)param;
783

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

      int64_t* pRefId = taosMemoryMalloc(sizeof(int64_t));
860
      *pRefId = pTmq->refId;
861

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

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

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

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

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

885
static void* tmqFreeRspWrapper(SMqRspWrapper* rspWrapper) {
L
Liu Jicong 已提交
886 887 888 889 890 891 892
  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;
893 894
    taosMemoryFreeClear(pRsp->pEpset);

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

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

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

  return NULL;
L
Liu Jicong 已提交
918 919
}

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

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

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

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

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

L
Liu Jicong 已提交
965
int32_t tmq_unsubscribe(tmq_t* tmq) {
L
Liu Jicong 已提交
966 967
  int32_t     rsp;
  int32_t     retryCnt = 0;
L
Liu Jicong 已提交
968
  tmq_list_t* lst = tmq_list_new();
L
Liu Jicong 已提交
969 970 971 972 973 974 975 976 977 978
  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 已提交
979 980
  tmq_list_destroy(lst);
  return rsp;
X
Xiaoyu Wang 已提交
981 982
}

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

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

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

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

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

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

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

1013 1014 1015 1016 1017 1018 1019 1020 1021
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);
1022
  if (tmqMgmt.rsetId < 0) {
1023 1024 1025 1026
    tmqInitRes = terrno;
  }
}

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

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

L
Liu Jicong 已提交
1041 1042 1043
  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 已提交
1044 1045 1046
  pTmq->clientTopics = taosArrayInit(0, sizeof(SMqClientTopic));
  pTmq->mqueue = taosOpenQueue();
  pTmq->delayedTask = taosOpenQueue();
H
Haojun Liao 已提交
1047
  pTmq->qall = taosAllocateQall();
L
Liu Jicong 已提交
1048

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

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

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

1072 1073
  pTmq->hbBgEnable = conf->hbBgEnable;

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

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

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

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

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

1103 1104 1105
  char         buf[80] = {0};
  STqOffsetVal offset = {.type = pTmq->resetOffsetCfg};
  tFormatOffset(buf, tListLen(buf), &offset);
X
Xiaoyu Wang 已提交
1106 1107 1108 1109
  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 已提交
1110

1111
  return pTmq;
1112

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

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

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

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

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

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

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

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

    taosArrayPush(req.topicNames, &topicFName);
1153 1154
  }

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

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

1163 1164 1165
  void* abuf = buf;
  tSerializeSCMSubscribeReq(&abuf, &req);

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

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

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

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

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

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

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

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

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

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

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

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

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

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

L
Liu Jicong 已提交
1234
  return code;
1235 1236
}

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

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

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

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

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

L
Liu Jicong 已提交
1262
  taosMemoryFree(pParam);
H
Haojun Liao 已提交
1263

L
Liu Jicong 已提交
1264
  if (code != 0) {
L
Liu Jicong 已提交
1265
    if (pMsg->pData) taosMemoryFree(pMsg->pData);
dengyihao's avatar
dengyihao 已提交
1266 1267
    if (pMsg->pEpSet) taosMemoryFree(pMsg->pEpSet);

H
Haojun Liao 已提交
1268
    // in case of consumer mismatch, wait for 500ms and retry
L
Liu Jicong 已提交
1269
    if (code == TSDB_CODE_TMQ_CONSUMER_MISMATCH) {
1270
      taosMsleep(500);
L
Liu Jicong 已提交
1271
      atomic_store_8(&tmq->status, TMQ_CONSUMER_STATUS__RECOVER);
X
Xiaoyu Wang 已提交
1272 1273
      tscDebug("consumer:0x%" PRIx64 " wait for the re-balance, wait for 500ms and set status to be RECOVER",
               tmq->consumerId);
H
Haojun Liao 已提交
1274
    } else if (code == TSDB_CODE_TQ_NO_COMMITTED_OFFSET) {
S
Shengliang Guan 已提交
1275
      SMqPollRspWrapper* pRspWrapper = taosAllocateQitem(sizeof(SMqPollRspWrapper), DEF_QITEM, 0);
L
Liu Jicong 已提交
1276
      if (pRspWrapper == NULL) {
H
Haojun Liao 已提交
1277 1278
        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 已提交
1279 1280
        goto CREATE_MSG_FAIL;
      }
H
Haojun Liao 已提交
1281

L
Liu Jicong 已提交
1282 1283
      pRspWrapper->tmqRspType = TMQ_MSG_TYPE__END_RSP;
      taosWriteQitem(tmq->mqueue, pRspWrapper);
X
Xiaoyu Wang 已提交
1284
    } else if (code == TSDB_CODE_WAL_LOG_NOT_EXIST) {  // poll data while insert
1285
      taosMsleep(500);
wmmhello's avatar
wmmhello 已提交
1286 1287 1288
    } else{
      tscError("consumer:0x%" PRIx64 " msg from vgId:%d discarded, epoch %d, since %s, reqId:0x%" PRIx64, tmq->consumerId,
               vgId, epoch, tstrerror(code), requestId);
L
Liu Jicong 已提交
1289
    }
H
Haojun Liao 已提交
1290

L
fix txn  
Liu Jicong 已提交
1291
    goto CREATE_MSG_FAIL;
L
Liu Jicong 已提交
1292 1293
  }

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

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

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

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

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

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

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

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

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

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

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

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

L
Liu Jicong 已提交
1372
  return 0;
H
Haojun Liao 已提交
1373

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

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

L
Liu Jicong 已提交
1382
  return -1;
1383 1384
}

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

H
Haojun Liao 已提交
1390 1391 1392 1393 1394 1395
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 已提交
1396
  char    vgKey[TSDB_TOPIC_FNAME_LEN + 22];
H
Haojun Liao 已提交
1397 1398 1399 1400 1401
  int32_t vgNumGet = taosArrayGetSize(pTopicEp->vgs);

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

wmmhello's avatar
wmmhello 已提交
1402
  tscDebug("consumer:0x%" PRIx64 ", update topic:%s, new numOfVgs:%d", tmq->consumerId, pTopic->topicName, vgNumGet);
H
Haojun Liao 已提交
1403 1404 1405 1406
  pTopic->vgs = taosArrayInit(vgNumGet, sizeof(SMqClientVg));

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

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

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

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

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

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

  taosArrayDestroy(pTopic->vgs);
}

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

1448
  int32_t topicNumCur = taosArrayGetSize(tmq->clientTopics);
1449
  int32_t topicNumGet = taosArrayGetSize(pRsp->topics);
1450

X
Xiaoyu Wang 已提交
1451 1452
  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",
1453
           tmq->consumerId, tmq->epoch, epoch, topicNumGet, topicNumCur);
wmmhello's avatar
wmmhello 已提交
1454 1455 1456
  if (epoch <= tmq->epoch) {
    return false;
  }
1457 1458 1459 1460 1461 1462

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

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

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

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

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

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

H
Haojun Liao 已提交
1498 1499
  taosHashCleanup(pVgOffsetHashMap);

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

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

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

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

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

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

H
Haojun Liao 已提交
1528
  if (code != TSDB_CODE_SUCCESS) {
1529 1530 1531 1532 1533 1534 1535 1536 1537
    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;
1538
  }
L
Liu Jicong 已提交
1539

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

1549 1550 1551 1552 1553 1554 1555 1556
    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 已提交
1557
  } else {
1558 1559
    tscDebug("consumer:0x%" PRIx64 ", recv ep, msg epoch %d, current epoch %d, update local ep", tmq->consumerId,
             head->epoch, epoch);
1560
  }
L
Liu Jicong 已提交
1561

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

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

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

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

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

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

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

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

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

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

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

L
Liu Jicong 已提交
1625
  return pRspObj;
X
Xiaoyu Wang 已提交
1626 1627
}

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

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

  return pRspObj;
}

1646 1647 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
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 已提交
1679
  pParam->pVg = pVg;
1680 1681
  pParam->pTopic = pTopic;
  pParam->vgId = pVg->vgId;
H
Haojun Liao 已提交
1682
  pParam->requestId = req.reqId;
1683 1684 1685 1686 1687 1688 1689 1690

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

H
Haojun Liao 已提交
1691
  sendInfo->msgInfo = (SDataBuf){.pData = msg, .len = msgSize, .handle = NULL};
1692 1693 1694 1695 1696 1697 1698 1699
  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 已提交
1700
  tFormatOffset(offsetFormatBuf, tListLen(offsetFormatBuf), &pVg->offsetInfo.currentOffset);
1701

X
Xiaoyu Wang 已提交
1702 1703
  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);
1704 1705 1706 1707 1708 1709 1710 1711
  asyncSendMsgToServer(pTmq->pTscObj->pAppInfo->pTransporter, &pVg->epSet, &transporterId, sendInfo);

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

  return TSDB_CODE_SUCCESS;
}

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

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

    for (int j = 0; j < numOfVg; j++) {
X
Xiaoyu Wang 已提交
1722
      SMqClientVg* pVg = taosArrayGet(pTopic->vgs, j);
X
Xiaoyu Wang 已提交
1723 1724 1725
      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 已提交
1726 1727 1728
        continue;
      }

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

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

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

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

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

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

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

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

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

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

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

        // update the epset
        if (pollRspWrapper->pEpset != NULL) {
          SEp* pEp = GET_ACTIVE_EP(pollRspWrapper->pEpset);
          SEp* pOld = GET_ACTIVE_EP(&(pVg->epSet));
X
Xiaoyu Wang 已提交
1804 1805
          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);
1806 1807 1808
          pVg->epSet = *pollRspWrapper->pEpset;
        }

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

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

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

1820 1821 1822
        char buf[80];
        tFormatOffset(buf, 80, &pDataRsp->rspOffset);
        if (pDataRsp->blockNum == 0) {
X
Xiaoyu Wang 已提交
1823 1824 1825
          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);
1826
          pRspWrapper = tmqFreeRspWrapper(pRspWrapper);
1827
          pVg->emptyBlockReceiveTs = taosGetTimestampMs();
L
Liu Jicong 已提交
1828
          taosFreeQitem(pollRspWrapper);
1829
        } else {  // build rsp
X
Xiaoyu Wang 已提交
1830
          int64_t    numOfRows = 0;
1831
          SMqRspObj* pRsp = tmqBuildRspFromWrapper(pollRspWrapper, pVg, &numOfRows);
H
Haojun Liao 已提交
1832
          tmq->totalRows += numOfRows;
1833
          pVg->emptyBlockReceiveTs = 0;
H
Haojun Liao 已提交
1834
          tscDebug("consumer:0x%" PRIx64 " process poll rsp, vgId:%d, offset:%s, blocks:%d, rows:%" PRId64
H
Haojun Liao 已提交
1835
                   " vg total:%" PRId64 " total:%" PRId64 ", reqId:0x%" PRIx64,
H
Haojun Liao 已提交
1836
                   tmq->consumerId, pVg->vgId, buf, pDataRsp->blockNum, numOfRows, pVg->numOfRows, tmq->totalRows,
H
Haojun Liao 已提交
1837
                   pollRspWrapper->reqId);
1838 1839 1840
          taosFreeQitem(pollRspWrapper);
          return pRsp;
        }
X
Xiaoyu Wang 已提交
1841
      } else {
H
Haojun Liao 已提交
1842
        tscDebug("consumer:0x%" PRIx64 " vgId:%d msg discard since epoch mismatch: msg epoch %d, consumer epoch %d",
1843
                 tmq->consumerId, pollRspWrapper->vgId, pDataRsp->head.epoch, consumerEpoch);
1844
        pRspWrapper = tmqFreeRspWrapper(pRspWrapper);
L
Liu Jicong 已提交
1845 1846
        taosFreeQitem(pollRspWrapper);
      }
1847
    } else if (pRspWrapper->tmqRspType == TMQ_MSG_TYPE__POLL_META_RSP) {
1848
      // todo handle the wal range and epset for each vgroup
1849
      SMqPollRspWrapper* pollRspWrapper = (SMqPollRspWrapper*)pRspWrapper;
L
Liu Jicong 已提交
1850
      int32_t            consumerEpoch = atomic_load_32(&tmq->epoch);
1851 1852 1853

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

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

L
Liu Jicong 已提交
1874 1875
      if (pollRspWrapper->taosxRsp.head.epoch == consumerEpoch) {
        SMqClientVg* pVg = pollRspWrapper->vgHandle;
1876
        if(pollRspWrapper->taosxRsp.rspOffset.type != 0){    // if offset is validate
H
Haojun Liao 已提交
1877
          pVg->offsetInfo.currentOffset = pollRspWrapper->taosxRsp.rspOffset;
1878
        }
L
Liu Jicong 已提交
1879
        atomic_store_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE);
H
Haojun Liao 已提交
1880

L
Liu Jicong 已提交
1881
        if (pollRspWrapper->taosxRsp.blockNum == 0) {
H
Haojun Liao 已提交
1882 1883
          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 已提交
1884
          pVg->emptyBlockReceiveTs = taosGetTimestampMs();
1885
          pRspWrapper = tmqFreeRspWrapper(pRspWrapper);
H
Haojun Liao 已提交
1886
          taosFreeQitem(pollRspWrapper);
L
Liu Jicong 已提交
1887
          continue;
H
Haojun Liao 已提交
1888
        } else {
X
Xiaoyu Wang 已提交
1889
          pVg->emptyBlockReceiveTs = 0;  // reset the ts
L
Liu Jicong 已提交
1890
        }
wmmhello's avatar
wmmhello 已提交
1891

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

1901 1902
        tmq->totalRows += numOfRows;

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

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

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

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

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

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

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

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

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

X
Xiaoyu Wang 已提交
1959
  while (1) {
L
Liu Jicong 已提交
1960
    tmqHandleAllDelayedTask(tmq);
1961

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

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

1975
    if (timeout >= 0) {
L
Liu Jicong 已提交
1976
      int64_t currentTime = taosGetTimestampMs();
1977 1978 1979
      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 已提交
1980
                 tmq->consumerId, tmq->epoch, startTime, currentTime);
X
Xiaoyu Wang 已提交
1981 1982
        return NULL;
      }
1983
      tsem_timewait(&tmq->rspSem, (timeout - elapsedTime));
L
Liu Jicong 已提交
1984 1985
    } else {
      // use tsem_timewait instead of tsem_wait to avoid unexpected stuck
L
Liu Jicong 已提交
1986
      tsem_timewait(&tmq->rspSem, 1000);
X
Xiaoyu Wang 已提交
1987 1988 1989 1990
    }
  }
}

1991 1992 1993 1994 1995 1996 1997 1998 1999 2000 2001 2002 2003 2004
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);
2005
    }
2006
  }
2007

2008 2009
  tscDebug("consumer:0x%" PRIx64 " rows dist end", pTmq->consumerId);
}
2010

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

2015 2016 2017 2018 2019 2020
  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;
2021 2022 2023
      }
    }

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

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

2041
  taosRemoveRef(tmqMgmt.rsetId, tmq->refId);
L
Liu Jicong 已提交
2042
  return 0;
2043
}
L
Liu Jicong 已提交
2044

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

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

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

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

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

2112 2113 2114 2115 2116 2117 2118 2119 2120 2121 2122 2123 2124 2125 2126 2127 2128 2129 2130 2131 2132 2133 2134
int64_t tmq_get_vgroup_offset(TAOS_RES* res) {
  if (TD_RES_TMQ(res)) {
    SMqRspObj* pRspObj = (SMqRspObj*) res;
    STqOffsetVal* pOffset = &pRspObj->rsp.rspOffset;
    if (pOffset->type == TMQ_OFFSET__LOG) {
      return pRspObj->rsp.rspOffset.version;
    }
  } else if (TD_RES_TMQ_META(res)) {
    SMqMetaRspObj* pRspObj = (SMqMetaRspObj*)res;
    if (pRspObj->metaRsp.rspOffset.type == TMQ_OFFSET__LOG) {
      return pRspObj->metaRsp.rspOffset.version;
    }
  } else if (TD_RES_TMQ_METADATA(res)) {
    SMqTaosxRspObj* pRspObj = (SMqTaosxRspObj*) res;
    if (pRspObj->rsp.rspOffset.type == TMQ_OFFSET__LOG) {
      return pRspObj->rsp.rspOffset.version;
    }
  }

  // data from tsdb, no valid offset info
  return -1;
}

L
Liu Jicong 已提交
2135 2136 2137 2138 2139 2140 2141
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;
    }
2142
    return (const char*)taosArrayGetP(pRspObj->rsp.blockTbName, pRspObj->resIter);
2143 2144
  } else if (TD_RES_TMQ_METADATA(res)) {
    SMqTaosxRspObj* pRspObj = (SMqTaosxRspObj*)res;
L
Liu Jicong 已提交
2145 2146 2147
    if (!pRspObj->rsp.withTbName || pRspObj->rsp.blockTbName == NULL || pRspObj->resIter < 0 ||
        pRspObj->resIter >= pRspObj->rsp.blockNum) {
      return NULL;
2148
    }
L
Liu Jicong 已提交
2149 2150
    return (const char*)taosArrayGetP(pRspObj->rsp.blockTbName, pRspObj->resIter);
  }
L
Liu Jicong 已提交
2151 2152
  return NULL;
}
2153

2154 2155 2156 2157
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
2158
    asyncCommitOffset(tmq, pRes, TDMT_VND_TMQ_COMMIT_OFFSET, cb, param);
2159
  }
L
Liu Jicong 已提交
2160 2161
}

2162
static void commitCallBackFn(tmq_t *UNUSED_PARAM(tmq), int32_t code, void* param) {
2163 2164 2165
  SSyncCommitInfo* pInfo = (SSyncCommitInfo*) param;
  pInfo->code = code;
  tsem_post(&pInfo->sem);
2166
}
2167

2168 2169 2170 2171 2172 2173 2174 2175 2176
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 已提交
2177
  } else {
2178
    asyncCommitOffset(tmq, pRes, TDMT_VND_TMQ_COMMIT_OFFSET, commitCallBackFn, pInfo);
2179 2180
  }

2181 2182
  tsem_wait(&pInfo->sem);
  code = pInfo->code;
H
Haojun Liao 已提交
2183 2184

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

X
Xiaoyu Wang 已提交
2187
  tscDebug("consumer:0x%" PRIx64 " sync commit done, code:%s", tmq->consumerId, tstrerror(code));
2188 2189 2190 2191 2192 2193 2194 2195 2196 2197 2198 2199
  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);
2200
    doUpdateLocalEp(pTmq, head->epoch, &rsp);
2201 2202 2203
    tDeleteSMqAskEpRsp(&rsp);
  }

H
Haojun Liao 已提交
2204
  tsem_post(&pInfo->sem);
2205 2206 2207 2208 2209 2210 2211 2212 2213 2214 2215 2216 2217 2218 2219 2220 2221 2222 2223 2224 2225 2226 2227 2228 2229 2230
}

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 已提交
2231
  tsem_init(&pInfo->sem, 0, 0);
2232 2233

  asyncAskEp(pTmq, updateEpCallbackFn, pInfo);
H
Haojun Liao 已提交
2234
  tsem_wait(&pInfo->sem);
2235 2236

  int32_t code = pInfo->code;
H
Haojun Liao 已提交
2237
  tsem_destroy(&pInfo->sem);
2238 2239 2240 2241 2242
  taosMemoryFree(pInfo);
  return code;
}

void asyncAskEp(tmq_t* pTmq, __tmq_askep_fn_t askEpFn, void* param) {
2243
  SMqAskEpReq req = {0};
2244 2245 2246
  req.consumerId = pTmq->consumerId;
  req.epoch = pTmq->epoch;
  strcpy(req.cgroup, pTmq->groupId);
2247 2248 2249

  int32_t tlen = tSerializeSMqAskEpReq(NULL, 0, &req);
  if (tlen < 0) {
2250 2251 2252
    tscError("consumer:0x%" PRIx64 ", tSerializeSMqAskEpReq failed", pTmq->consumerId);
    askEpFn(pTmq, TSDB_CODE_INVALID_PARA, NULL, param);
    return;
2253 2254 2255 2256
  }

  void* pReq = taosMemoryCalloc(1, tlen);
  if (pReq == NULL) {
2257 2258 2259
    tscError("consumer:0x%" PRIx64 ", failed to malloc askEpReq msg, size:%d", pTmq->consumerId, tlen);
    askEpFn(pTmq, TSDB_CODE_OUT_OF_MEMORY, NULL, param);
    return;
2260 2261 2262
  }

  if (tSerializeSMqAskEpReq(pReq, tlen, &req) < 0) {
2263
    tscError("consumer:0x%" PRIx64 ", tSerializeSMqAskEpReq %d failed", pTmq->consumerId, tlen);
2264
    taosMemoryFree(pReq);
2265 2266 2267

    askEpFn(pTmq, TSDB_CODE_INVALID_PARA, NULL, param);
    return;
2268 2269 2270 2271
  }

  SMqAskEpCbParam* pParam = taosMemoryCalloc(1, sizeof(SMqAskEpCbParam));
  if (pParam == NULL) {
2272
    tscError("consumer:0x%" PRIx64 ", failed to malloc subscribe param", pTmq->consumerId);
2273
    taosMemoryFree(pReq);
2274 2275 2276

    askEpFn(pTmq, TSDB_CODE_OUT_OF_MEMORY, NULL, param);
    return;
2277 2278
  }

2279 2280 2281 2282
  pParam->refId = pTmq->refId;
  pParam->epoch = pTmq->epoch;
  pParam->pUserFn = askEpFn;
  pParam->pParam = param;
2283 2284 2285 2286 2287

  SMsgSendInfo* sendInfo = taosMemoryCalloc(1, sizeof(SMsgSendInfo));
  if (sendInfo == NULL) {
    taosMemoryFree(pParam);
    taosMemoryFree(pReq);
2288 2289
    askEpFn(pTmq, TSDB_CODE_OUT_OF_MEMORY, NULL, param);
    return;
2290 2291
  }

X
Xiaoyu Wang 已提交
2292
  sendInfo->msgInfo = (SDataBuf){.pData = pReq, .len = tlen, .handle = NULL};
2293 2294 2295 2296

  sendInfo->requestId = generateRequestId();
  sendInfo->requestObjRefId = 0;
  sendInfo->param = pParam;
2297
  sendInfo->fp = askEpCallbackFn;
2298 2299
  sendInfo->msgType = TDMT_MND_TMQ_ASK_EP;

2300 2301
  SEpSet epSet = getEpSet_s(&pTmq->pTscObj->pAppInfo->mgmtEp);
  tscDebug("consumer:0x%" PRIx64 " ask ep from mnode, reqId:0x%" PRIx64, pTmq->consumerId, sendInfo->requestId);
2302 2303

  int64_t transporterId = 0;
2304
  asyncSendMsgToServer(pTmq->pTscObj->pAppInfo->pTransporter, &epSet, &transporterId, sendInfo);
2305 2306 2307 2308 2309 2310 2311
}

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

int32_t tmqCommitDone(SMqCommitCbParamSet* pParamSet) {
2312 2313 2314
  int64_t refId = pParamSet->refId;

  tmq_t* tmq = taosAcquireRef(tmqMgmt.rsetId, refId);
2315 2316 2317 2318 2319 2320 2321
  if (tmq == NULL) {
    taosMemoryFree(pParamSet);
    terrno = TSDB_CODE_TMQ_CONSUMER_CLOSED;
    return -1;
  }

  // if no more waiting rsp
2322
  pParamSet->callbackFn(tmq, pParamSet->code, pParamSet->userParam);
2323
  taosMemoryFree(pParamSet);
2324 2325

  taosReleaseRef(tmqMgmt.rsetId, refId);
2326
  return 0;
2327 2328
}

2329
void commitRspCountDown(SMqCommitCbParamSet* pParamSet, int64_t consumerId, const char* pTopic, int32_t vgId) {
2330 2331
  int32_t waitingRspNum = atomic_sub_fetch_32(&pParamSet->waitingRspNum, 1);
  if (waitingRspNum == 0) {
H
Haojun Liao 已提交
2332 2333
    tscDebug("consumer:0x%" PRIx64 " topic:%s vgId:%d all commit-rsp received, commit completed", consumerId, pTopic,
             vgId);
2334
    tmqCommitDone(pParamSet);
H
Haojun Liao 已提交
2335 2336 2337
  } else {
    tscDebug("consumer:0x%" PRIx64 " topic:%s vgId:%d commit-rsp received, remain:%d", consumerId, pTopic, vgId,
             waitingRspNum);
2338 2339
  }
}
2340 2341 2342 2343 2344 2345 2346 2347 2348 2349 2350 2351 2352 2353 2354 2355 2356 2357 2358 2359 2360 2361

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 已提交
2362 2363
}

2364 2365 2366 2367 2368 2369 2370 2371 2372 2373 2374 2375 2376 2377 2378 2379 2380 2381 2382 2383 2384
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,
2385
                                       .vgId = pParam->vgId};
2386 2387 2388 2389 2390 2391 2392 2393 2394 2395 2396 2397 2398 2399 2400 2401 2402 2403 2404 2405 2406 2407

    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 已提交
2408
int32_t tmq_get_topic_assignment(tmq_t* tmq, const char* pTopicName, tmq_topic_assignment** assignment,
H
Haojun Liao 已提交
2409
                                 int32_t* numOfAssignment) {
H
Haojun Liao 已提交
2410 2411 2412
  *numOfAssignment = 0;
  *assignment = NULL;

2413
  int32_t accId = tmq->pTscObj->acctId;
2414
  char    tname[128] = {0};
2415 2416 2417
  sprintf(tname, "%d.%s", accId, pTopicName);

  SMqClientTopic* pTopic = getTopicByName(tmq, tname);
H
Haojun Liao 已提交
2418 2419 2420 2421 2422 2423
  if (pTopic == NULL) {
    return TSDB_CODE_INVALID_PARA;
  }

  // in case of snapshot is opened, no valid offset will return
  *numOfAssignment = taosArrayGetSize(pTopic->vgs);
2424 2425 2426 2427 2428 2429 2430 2431

  *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;
  }

2432 2433
  bool needFetch = false;

H
Haojun Liao 已提交
2434 2435
  for (int32_t j = 0; j < (*numOfAssignment); ++j) {
    SMqClientVg* pClientVg = taosArrayGet(pTopic->vgs, j);
2436 2437 2438 2439
    if (!pClientVg->receiveInfo) {
      needFetch = true;
      break;
    }
H
Haojun Liao 已提交
2440 2441 2442 2443 2444 2445 2446 2447 2448 2449

    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;
2450
    pAssignment->vgId = pClientVg->vgId;
H
Haojun Liao 已提交
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 2523 2524 2525 2526 2527 2528 2529 2530 2531 2532 2533 2534
  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);
2535
      *assignment = NULL;
2536 2537 2538 2539 2540 2541 2542 2543 2544
      *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;
    }

2545 2546 2547 2548 2549 2550 2551 2552 2553 2554 2555 2556 2557 2558 2559 2560 2561 2562 2563 2564 2565 2566 2567 2568 2569
    for (int32_t j = 0; j < (*numOfAssignment); ++j) {
      tmq_topic_assignment* p = &(*assignment)[j];

      for(int32_t i = 0; i < taosArrayGetSize(pTopic->vgs); ++i) {
        SMqClientVg* pClientVg = taosArrayGet(pTopic->vgs, i);
        if (pClientVg->vgId != p->vgId) {
          continue;
        }

        SVgOffsetInfo* pOffsetInfo = &pClientVg->offsetInfo;

        pOffsetInfo->currentOffset.type = TMQ_OFFSET__LOG;

        char offsetBuf[80] = {0};
        tFormatOffset(offsetBuf, tListLen(offsetBuf), &pOffsetInfo->currentOffset);

        tscDebug("vgId:%d offset is update to:%s", p->vgId, offsetBuf);

        pOffsetInfo->walVerBegin = p->begin;
        pOffsetInfo->walVerEnd = p->end;
        pOffsetInfo->currentOffset.version = p->currentOffset;
        pOffsetInfo->committedOffset.version = p->currentOffset;
      }
    }

2570 2571 2572 2573 2574
    destroyCommonInfo(pCommon);
    return code;
  } else {
    return TSDB_CODE_SUCCESS;
  }
H
Haojun Liao 已提交
2575 2576
}

2577
int32_t tmq_offset_seek(tmq_t* tmq, const char* pTopicName, int32_t vgId, int64_t offset) {
H
Haojun Liao 已提交
2578
  if (tmq == NULL) {
H
Haojun Liao 已提交
2579
    tscError("invalid tmq handle, null");
H
Haojun Liao 已提交
2580 2581 2582
    return TSDB_CODE_INVALID_PARA;
  }

2583 2584 2585 2586 2587
  int32_t accId = tmq->pTscObj->acctId;
  char tname[128] = {0};
  sprintf(tname, "%d.%s", accId, pTopicName);

  SMqClientTopic* pTopic = getTopicByName(tmq, tname);
H
Haojun Liao 已提交
2588
  if (pTopic == NULL) {
2589
    tscError("consumer:0x%" PRIx64 " invalid topic name:%s", tmq->consumerId, pTopicName);
H
Haojun Liao 已提交
2590 2591 2592 2593
    return TSDB_CODE_INVALID_PARA;
  }

  SMqClientVg* pVg = NULL;
H
Haojun Liao 已提交
2594 2595
  int32_t      numOfVgs = taosArrayGetSize(pTopic->vgs);
  for (int32_t i = 0; i < numOfVgs; ++i) {
H
Haojun Liao 已提交
2596
    SMqClientVg* pClientVg = taosArrayGet(pTopic->vgs, i);
2597
    if (pClientVg->vgId == vgId) {
H
Haojun Liao 已提交
2598 2599 2600 2601 2602 2603
      pVg = pClientVg;
      break;
    }
  }

  if (pVg == NULL) {
2604
    tscError("consumer:0x%" PRIx64 " invalid vgroup id:%d", tmq->consumerId, vgId);
H
Haojun Liao 已提交
2605 2606 2607 2608 2609 2610 2611
    return TSDB_CODE_INVALID_PARA;
  }

  SVgOffsetInfo* pOffsetInfo = &pVg->offsetInfo;

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

H
Haojun Liao 已提交
2616
  if (offset < pOffsetInfo->walVerBegin || offset > pOffsetInfo->walVerEnd) {
H
Haojun Liao 已提交
2617 2618
    tscError("consumer:0x%" PRIx64 " invalid seek params, offset:%" PRId64 ", valid range:[%" PRId64 ", %" PRId64 "]",
             tmq->consumerId, offset, pOffsetInfo->walVerBegin, pOffsetInfo->walVerEnd);
H
Haojun Liao 已提交
2619 2620 2621
    return TSDB_CODE_INVALID_PARA;
  }

H
Haojun Liao 已提交
2622 2623 2624
  // update the offset, and then commit to vnode
  if (pOffsetInfo->currentOffset.type == TMQ_OFFSET__LOG) {
    pOffsetInfo->currentOffset.version = offset;
2625
    pOffsetInfo->committedOffset.version = INT64_MIN;
H
Haojun Liao 已提交
2626 2627
  }

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

H
Haojun Liao 已提交
2631
  tscDebug("consumer:0x%" PRIx64 " seek to %" PRId64 " on vgId:%d", tmq->consumerId, offset, pVg->vgId);
2632 2633 2634 2635 2636 2637 2638 2639 2640 2641 2642 2643 2644 2645 2646 2647 2648 2649

  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 已提交
2650 2651 2652 2653 2654 2655
  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;
2656
}