clientTmq.c 86.1 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
#define OFFSET_IS_RESET_OFFSET(_of)  ((_of) < 0)

32 33
typedef void (*__tmq_askep_fn_t)(tmq_t* pTmq, int32_t code, SDataBuf* pBuf, void* pParam);

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

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

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

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

L
Liu Jicong 已提交
55
struct tmq_list_t {
L
Liu Jicong 已提交
56
  SArray container;
L
Liu Jicong 已提交
57
};
L
Liu Jicong 已提交
58

L
Liu Jicong 已提交
59
struct tmq_conf_t {
60 61 62 63 64 65 66 67
  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;
68 69 70 71 72
  uint16_t       port;
  int32_t        autoCommitInterval;
  char*          ip;
  char*          user;
  char*          pass;
73
  tmq_commit_cb* commitCb;
L
Liu Jicong 已提交
74
  void*          commitCbUserParam;
L
Liu Jicong 已提交
75 76 77
};

struct tmq_t {
78 79 80 81 82 83 84 85 86 87
  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 已提交
88 89
  tmq_commit_cb* commitCb;
  void*          commitCbUserParam;
L
Liu Jicong 已提交
90 91 92 93

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

L
Liu Jicong 已提交
102
  // timer
X
Xiaoyu Wang 已提交
103 104 105 106 107 108 109 110 111 112
  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 已提交
113 114
};

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

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

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

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

H
Haojun Liao 已提交
138
typedef struct SVgOffsetInfo {
L
Liu Jicong 已提交
139 140
  STqOffsetVal committedOffset;
  STqOffsetVal currentOffset;
H
Haojun Liao 已提交
141 142 143 144 145 146 147 148 149 150 151
  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
152
  bool          receivedInfoFromVnode;// has already received info from vnode
H
Haojun Liao 已提交
153
  int64_t       emptyBlockReceiveTs;  // once empty block is received, idle for ignoreCnt then start to poll data
154
  bool          seekUpdated;          // offset is updated by seek operator, therefore, not update by vnode rsp.
H
Haojun Liao 已提交
155
  SEpSet        epSet;
156 157
} SMqClientVg;

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

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

L
Liu Jicong 已提交
180
typedef struct {
181 182
  int64_t refId;
  int32_t epoch;
L
Liu Jicong 已提交
183 184
  tsem_t  rspSem;
  int32_t rspErr;
L
Liu Jicong 已提交
185
} SMqSubscribeCbParam;
L
Liu Jicong 已提交
186

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

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

203 204 205 206 207 208 209 210 211 212 213 214 215 216 217 218 219
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;

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

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

240 241 242 243 244
typedef struct SSyncCommitInfo {
  tsem_t  sem;
  int32_t code;
} SSyncCommitInfo;

245
static int32_t doAskEp(tmq_t* tmq);
246 247
static int32_t makeTopicVgroupKey(char* dst, const char* topicName, int32_t vg);
static int32_t tmqCommitDone(SMqCommitCbParamSet* pParamSet);
248
static int32_t doSendCommitMsg(tmq_t* tmq, SMqClientVg* pVg, const char* pTopicName, SMqCommitCbParamSet* pParamSet,
249
                               int32_t index, int32_t totalVgroups, int32_t type);
X
Xiaoyu Wang 已提交
250 251 252
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);
253

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

261
  conf->withTbName = false;
L
Liu Jicong 已提交
262
  conf->autoCommit = true;
263
  conf->autoCommitInterval = DEFAULT_AUTO_COMMIT_INTERVAL;
264
  conf->resetOffset = TMQ_OFFSET__RESET_EARLIEAST;
265
  conf->hbBgEnable = true;
266

267 268 269
  return conf;
}

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

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

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

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

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

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

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

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

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

357
  if (strcasecmp(key, "enable.heartbeat.background") == 0) {
X
Xiaoyu Wang 已提交
358 359 360 361 362 363 364 365 366 367
    //    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 已提交
368 369
  }

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

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

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

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

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

L
Liu Jicong 已提交
394
  return TMQ_CONF_UNKNOWN;
395 396
}

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

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

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

L
Liu Jicong 已提交
412 413 414 415 416 417 418 419 420 421
int32_t tmq_list_get_size(const tmq_list_t* list) {
  const SArray* container = &list->container;
  return taosArrayGetSize(container);
}

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

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

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

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

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

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

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

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

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

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

500 501
  pOffset->consumerId = tmq->consumerId;
  pOffset->offset.val = pVg->offsetInfo.currentOffset;
502

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

648 649 650 651 652 653
  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 已提交
654
    }
L
Liu Jicong 已提交
655
  }
L
Liu Jicong 已提交
656

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

1114
  return pTmq;
1115

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

L
Liu Jicong 已提交
1285 1286
      pRspWrapper->tmqRspType = TMQ_MSG_TYPE__END_RSP;
      taosWriteQitem(tmq->mqueue, pRspWrapper);
X
Xiaoyu Wang 已提交
1287
    } else if (code == TSDB_CODE_WAL_LOG_NOT_EXIST) {  // poll data while insert
1288
      taosMsleep(500);
wmmhello's avatar
wmmhello 已提交
1289 1290 1291
    } 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 已提交
1292
    }
H
Haojun Liao 已提交
1293

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

H
Haojun Liao 已提交
1393 1394 1395 1396 1397 1398
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 已提交
1399
  char    vgKey[TSDB_TOPIC_FNAME_LEN + 22];
H
Haojun Liao 已提交
1400 1401 1402 1403 1404
  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 已提交
1405
  tscDebug("consumer:0x%" PRIx64 ", update topic:%s, new numOfVgs:%d", tmq->consumerId, pTopic->topicName, vgNumGet);
H
Haojun Liao 已提交
1406 1407 1408 1409
  pTopic->vgs = taosArrayInit(vgNumGet, sizeof(SMqClientVg));

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

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

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

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

H
Haojun Liao 已提交
1431 1432 1433 1434
    clientVg.offsetInfo.currentOffset = offsetNew;
    clientVg.offsetInfo.committedOffset = offsetNew;
    clientVg.offsetInfo.walVerBegin = -1;
    clientVg.offsetInfo.walVerEnd = -1;
1435 1436 1437
    clientVg.seekUpdated = false;
    clientVg.receivedInfoFromVnode = false;

H
Haojun Liao 已提交
1438 1439 1440 1441 1442 1443 1444 1445 1446 1447 1448 1449 1450
    taosArrayPush(pTopic->vgs, &clientVg);
  }
}

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

  taosArrayDestroy(pTopic->vgs);
}

1451
static bool doUpdateLocalEp(tmq_t* tmq, int32_t epoch, const SMqAskEpRsp* pRsp) {
1452 1453
  bool set = false;

1454
  int32_t topicNumCur = taosArrayGetSize(tmq->clientTopics);
1455
  int32_t topicNumGet = taosArrayGetSize(pRsp->topics);
1456

X
Xiaoyu Wang 已提交
1457 1458
  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",
1459
           tmq->consumerId, tmq->epoch, epoch, topicNumGet, topicNumCur);
wmmhello's avatar
wmmhello 已提交
1460 1461 1462
  if (epoch <= tmq->epoch) {
    return false;
  }
1463 1464 1465 1466 1467 1468

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

H
Haojun Liao 已提交
1469 1470
  SHashObj* pVgOffsetHashMap = taosHashInit(64, MurmurHash3_32, false, HASH_NO_LOCK);
  if (pVgOffsetHashMap == NULL) {
1471 1472 1473
    taosArrayDestroy(newTopics);
    return false;
  }
1474

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

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

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

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

H
Haojun Liao 已提交
1504 1505
  taosHashCleanup(pVgOffsetHashMap);

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

X
Xiaoyu Wang 已提交
1512
  int8_t flag = (topicNumGet == 0) ? TMQ_CONSUMER_STATUS__NO_TOPIC : TMQ_CONSUMER_STATUS__READY;
H
Haojun Liao 已提交
1513
  atomic_store_8(&tmq->status, flag);
X
Xiaoyu Wang 已提交
1514
  atomic_store_32(&tmq->epoch, epoch);
H
Haojun Liao 已提交
1515

1516
  tscDebug("consumer:0x%" PRIx64 " update topic info completed", tmq->consumerId);
X
Xiaoyu Wang 已提交
1517 1518 1519
  return set;
}

1520
int32_t askEpCallbackFn(void* param, SDataBuf* pMsg, int32_t code) {
1521
  SMqAskEpCbParam* pParam = (SMqAskEpCbParam*)param;
1522 1523 1524
  tmq_t*           tmq = taosAcquireRef(tmqMgmt.rsetId, pParam->refId);

  if (tmq == NULL) {
1525 1526 1527
    terrno = TSDB_CODE_TMQ_CONSUMER_CLOSED;
    pParam->pUserFn(tmq, terrno, NULL, pParam->pParam);

1528
    taosMemoryFree(pMsg->pData);
dengyihao's avatar
dengyihao 已提交
1529
    taosMemoryFree(pMsg->pEpSet);
1530 1531
    taosMemoryFree(pParam);
    return terrno;
1532 1533
  }

H
Haojun Liao 已提交
1534
  if (code != TSDB_CODE_SUCCESS) {
1535 1536 1537 1538 1539 1540 1541 1542 1543
    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;
1544
  }
L
Liu Jicong 已提交
1545

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

1555 1556 1557 1558 1559 1560 1561 1562
    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 已提交
1563
  } else {
1564 1565
    tscDebug("consumer:0x%" PRIx64 ", recv ep, msg epoch %d, current epoch %d, update local ep", tmq->consumerId,
             head->epoch, epoch);
1566
  }
L
Liu Jicong 已提交
1567

1568
  pParam->pUserFn(tmq, code, pMsg, pParam->pParam);
1569 1570
  taosReleaseRef(tmqMgmt.rsetId, pParam->refId);

dengyihao's avatar
dengyihao 已提交
1571
  taosMemoryFree(pMsg->pEpSet);
L
Liu Jicong 已提交
1572
  taosMemoryFree(pMsg->pData);
1573
  taosMemoryFree(pParam);
L
Liu Jicong 已提交
1574
  return code;
1575 1576
}

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

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

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

1604
SMqRspObj* tmqBuildRspFromWrapper(SMqPollRspWrapper* pWrapper, SMqClientVg* pVg, int64_t* numOfRows) {
L
Liu Jicong 已提交
1605 1606
  SMqRspObj* pRspObj = taosMemoryCalloc(1, sizeof(SMqRspObj));
  pRspObj->resType = RES_TYPE__TMQ;
1607

1608
  (*numOfRows) = 0;
L
Liu Jicong 已提交
1609 1610
  tstrncpy(pRspObj->topic, pWrapper->topicHandle->topicName, TSDB_TOPIC_FNAME_LEN);
  tstrncpy(pRspObj->db, pWrapper->topicHandle->db, TSDB_DB_FNAME_LEN);
1611

L
Liu Jicong 已提交
1612
  pRspObj->vgId = pWrapper->vgHandle->vgId;
L
Liu Jicong 已提交
1613
  pRspObj->resIter = -1;
L
Liu Jicong 已提交
1614
  memcpy(&pRspObj->rsp, &pWrapper->dataRsp, sizeof(SMqDataRsp));
L
Liu Jicong 已提交
1615

L
Liu Jicong 已提交
1616 1617
  pRspObj->resInfo.totalRows = 0;
  pRspObj->resInfo.precision = TSDB_TIME_PRECISION_MILLI;
1618

L
Liu Jicong 已提交
1619
  if (!pWrapper->dataRsp.withSchema) {
L
Liu Jicong 已提交
1620 1621
    setResSchemaInfo(&pRspObj->resInfo, pWrapper->topicHandle->schema.pSchema, pWrapper->topicHandle->schema.nCols);
  }
L
Liu Jicong 已提交
1622

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

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

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

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

  return pRspObj;
}

1652 1653 1654 1655 1656 1657 1658 1659 1660 1661 1662 1663 1664 1665 1666 1667 1668 1669 1670 1671 1672 1673 1674 1675 1676 1677 1678 1679 1680 1681 1682 1683 1684
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 已提交
1685
  pParam->pVg = pVg;
1686 1687
  pParam->pTopic = pTopic;
  pParam->vgId = pVg->vgId;
H
Haojun Liao 已提交
1688
  pParam->requestId = req.reqId;
1689 1690 1691 1692 1693 1694 1695 1696

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

H
Haojun Liao 已提交
1697
  sendInfo->msgInfo = (SDataBuf){.pData = msg, .len = msgSize, .handle = NULL};
1698 1699 1700 1701 1702 1703 1704 1705
  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 已提交
1706
  tFormatOffset(offsetFormatBuf, tListLen(offsetFormatBuf), &pVg->offsetInfo.currentOffset);
1707

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

  pVg->pollCnt++;
1713
  pVg->seekUpdated = false;   // reset this flag.
1714 1715 1716 1717 1718
  pTmq->pollCnt++;

  return TSDB_CODE_SUCCESS;
}

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

  for (int i = 0; i < numOfTopics; i++) {
X
Xiaoyu Wang 已提交
1725
    SMqClientTopic* pTopic = taosArrayGet(tmq->clientTopics, i);
X
Xiaoyu Wang 已提交
1726
    int32_t         numOfVg = taosArrayGetSize(pTopic->vgs);
1727 1728

    for (int j = 0; j < numOfVg; j++) {
X
Xiaoyu Wang 已提交
1729
      SMqClientVg* pVg = taosArrayGet(pTopic->vgs, j);
X
Xiaoyu Wang 已提交
1730 1731 1732
      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 已提交
1733 1734 1735
        continue;
      }

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

L
Liu Jicong 已提交
1744
      atomic_store_32(&pVg->vgSkipCnt, 0);
1745 1746 1747
      int32_t code = doTmqPollImpl(tmq, pTopic, pVg, timeout);
      if (code != TSDB_CODE_SUCCESS) {
        return code;
D
dapan1121 已提交
1748
      }
X
Xiaoyu Wang 已提交
1749 1750
    }
  }
1751

1752
  tscDebug("consumer:0x%" PRIx64 " end to poll data", tmq->consumerId);
X
Xiaoyu Wang 已提交
1753 1754 1755
  return 0;
}

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

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

X
Xiaoyu Wang 已提交
1779
  while (1) {
1780 1781
    SMqRspWrapper* pRspWrapper = NULL;
    taosGetQitem(tmq->qall, (void**)&pRspWrapper);
1782

1783
    if (pRspWrapper == NULL) {
L
Liu Jicong 已提交
1784
      taosReadAllQitems(tmq->mqueue, tmq->qall);
1785 1786
      taosGetQitem(tmq->qall, (void**)&pRspWrapper);
      if (pRspWrapper == NULL) {
L
Liu Jicong 已提交
1787 1788
        return NULL;
      }
X
Xiaoyu Wang 已提交
1789 1790
    }

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

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

X
Xiaoyu Wang 已提交
1801
      int32_t     consumerEpoch = atomic_load_32(&tmq->epoch);
1802 1803 1804
      SMqDataRsp* pDataRsp = &pollRspWrapper->dataRsp;

      if (pDataRsp->head.epoch == consumerEpoch) {
L
Liu Jicong 已提交
1805
        SMqClientVg* pVg = pollRspWrapper->vgHandle;
1806 1807 1808 1809 1810

        // update the epset
        if (pollRspWrapper->pEpset != NULL) {
          SEp* pEp = GET_ACTIVE_EP(pollRspWrapper->pEpset);
          SEp* pOld = GET_ACTIVE_EP(&(pVg->epSet));
X
Xiaoyu Wang 已提交
1811 1812
          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);
1813 1814 1815
          pVg->epSet = *pollRspWrapper->pEpset;
        }

1816 1817 1818 1819 1820
        // update the local offset value only for the returned values, only when the local offset is NOT updated
        // by tmq_offset_seek function
        if (!pVg->seekUpdated) {
          pVg->offsetInfo.currentOffset = pDataRsp->rspOffset;
        }
1821 1822

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

1825 1826 1827
        // update the valid wal version range
        pVg->offsetInfo.walVerBegin = pDataRsp->head.walsver;
        pVg->offsetInfo.walVerEnd = pDataRsp->head.walever;
1828
        pVg->receivedInfoFromVnode = true;
1829

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

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

L
Liu Jicong 已提交
1864
      if (pollRspWrapper->metaRsp.head.epoch == consumerEpoch) {
L
Liu Jicong 已提交
1865
        SMqClientVg* pVg = pollRspWrapper->vgHandle;
1866
        if (!pVg->seekUpdated) {
H
Haojun Liao 已提交
1867
          pVg->offsetInfo.currentOffset = pollRspWrapper->metaRsp.rspOffset;
1868
        }
1869

L
Liu Jicong 已提交
1870 1871
        atomic_store_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE);
        // build rsp
L
Liu Jicong 已提交
1872
        SMqMetaRspObj* pRsp = tmqBuildMetaRspFromWrapper(pollRspWrapper);
L
Liu Jicong 已提交
1873 1874 1875
        taosFreeQitem(pollRspWrapper);
        return pRsp;
      } else {
H
Haojun Liao 已提交
1876
        tscDebug("consumer:0x%" PRIx64 " vgId:%d msg discard since epoch mismatch: msg epoch %d, consumer epoch %d",
1877
                 tmq->consumerId, pollRspWrapper->vgId, pollRspWrapper->metaRsp.head.epoch, consumerEpoch);
1878
        pRspWrapper = tmqFreeRspWrapper(pRspWrapper);
L
Liu Jicong 已提交
1879
        taosFreeQitem(pollRspWrapper);
X
Xiaoyu Wang 已提交
1880
      }
1881 1882
    } else if (pRspWrapper->tmqRspType == TMQ_MSG_TYPE__TAOSX_RSP) {
      SMqPollRspWrapper* pollRspWrapper = (SMqPollRspWrapper*)pRspWrapper;
X
Xiaoyu Wang 已提交
1883
      int32_t            consumerEpoch = atomic_load_32(&tmq->epoch);
H
Haojun Liao 已提交
1884

L
Liu Jicong 已提交
1885 1886
      if (pollRspWrapper->taosxRsp.head.epoch == consumerEpoch) {
        SMqClientVg* pVg = pollRspWrapper->vgHandle;
1887
        if (!pVg->seekUpdated) {  // if offset is validate
H
Haojun Liao 已提交
1888
          pVg->offsetInfo.currentOffset = pollRspWrapper->taosxRsp.rspOffset;
1889
        }
1890

L
Liu Jicong 已提交
1891
        atomic_store_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE);
H
Haojun Liao 已提交
1892

L
Liu Jicong 已提交
1893
        if (pollRspWrapper->taosxRsp.blockNum == 0) {
H
Haojun Liao 已提交
1894 1895
          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 已提交
1896
          pVg->emptyBlockReceiveTs = taosGetTimestampMs();
1897
          pRspWrapper = tmqFreeRspWrapper(pRspWrapper);
H
Haojun Liao 已提交
1898
          taosFreeQitem(pollRspWrapper);
L
Liu Jicong 已提交
1899
          continue;
H
Haojun Liao 已提交
1900
        } else {
X
Xiaoyu Wang 已提交
1901
          pVg->emptyBlockReceiveTs = 0;  // reset the ts
L
Liu Jicong 已提交
1902
        }
wmmhello's avatar
wmmhello 已提交
1903

L
Liu Jicong 已提交
1904
        // build rsp
X
Xiaoyu Wang 已提交
1905
        void*   pRsp = NULL;
1906
        int64_t numOfRows = 0;
L
Liu Jicong 已提交
1907
        if (pollRspWrapper->taosxRsp.createTableNum == 0) {
1908
          pRsp = tmqBuildRspFromWrapper(pollRspWrapper, pVg, &numOfRows);
L
Liu Jicong 已提交
1909
        } else {
wmmhello's avatar
wmmhello 已提交
1910 1911
          pRsp = tmqBuildTaosxRspFromWrapper(pollRspWrapper);
        }
H
Haojun Liao 已提交
1912

1913 1914
        tmq->totalRows += numOfRows;

H
Haojun Liao 已提交
1915
        char buf[80];
H
Haojun Liao 已提交
1916
        tFormatOffset(buf, 80, &pVg->offsetInfo.currentOffset);
H
Haojun Liao 已提交
1917
        tscDebug("consumer:0x%" PRIx64 " process taosx poll rsp, vgId:%d, offset:%s, blocks:%d, rows:%" PRId64
X
Xiaoyu Wang 已提交
1918
                 ", vg total:%" PRId64 " total:%" PRId64 " reqId:0x%" PRIx64,
H
Haojun Liao 已提交
1919
                 tmq->consumerId, pVg->vgId, buf, pollRspWrapper->dataRsp.blockNum, numOfRows, pVg->numOfRows,
H
Haojun Liao 已提交
1920
                 tmq->totalRows, pollRspWrapper->reqId);
H
Haojun Liao 已提交
1921 1922

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

L
Liu Jicong 已提交
1925
      } else {
H
Haojun Liao 已提交
1926
        tscDebug("consumer:0x%" PRIx64 " vgId:%d msg discard since epoch mismatch: msg epoch %d, consumer epoch %d",
1927
                 tmq->consumerId, pollRspWrapper->vgId, pollRspWrapper->taosxRsp.head.epoch, consumerEpoch);
1928
        pRspWrapper = tmqFreeRspWrapper(pRspWrapper);
L
Liu Jicong 已提交
1929 1930
        taosFreeQitem(pollRspWrapper);
      }
X
Xiaoyu Wang 已提交
1931
    } else {
H
Haojun Liao 已提交
1932 1933
      tscDebug("consumer:0x%" PRIx64 " not data msg received", tmq->consumerId);

X
Xiaoyu Wang 已提交
1934
      bool reset = false;
1935 1936
      tmqHandleNoPollRsp(tmq, pRspWrapper, &reset);
      taosFreeQitem(pRspWrapper);
X
Xiaoyu Wang 已提交
1937
      if (pollIfReset && reset) {
1938
        tscDebug("consumer:0x%" PRIx64 ", reset and repoll", tmq->consumerId);
1939
        tmqPollImpl(tmq, timeout);
X
Xiaoyu Wang 已提交
1940 1941 1942 1943 1944
      }
    }
  }
}

1945
TAOS_RES* tmq_consumer_poll(tmq_t* tmq, int64_t timeout) {
L
Liu Jicong 已提交
1946 1947
  void*   rspObj;
  int64_t startTime = taosGetTimestampMs();
L
Liu Jicong 已提交
1948

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

1952
  // in no topic status, delayed task also need to be processed
L
Liu Jicong 已提交
1953
  if (atomic_load_8(&tmq->status) == TMQ_CONSUMER_STATUS__INIT) {
1954
    tscDebug("consumer:0x%" PRIx64 " poll return since consumer is init", tmq->consumerId);
1955
    taosMsleep(500);  //     sleep for a while
1956 1957 1958
    return NULL;
  }

wmmhello's avatar
wmmhello 已提交
1959
  while (atomic_load_8(&tmq->status) == TMQ_CONSUMER_STATUS__RECOVER) {
L
Liu Jicong 已提交
1960
    int32_t retryCnt = 0;
1961
    while (TSDB_CODE_MND_CONSUMER_NOT_READY == doAskEp(tmq)) {
H
Haojun Liao 已提交
1962
      if (retryCnt++ > 40) {
L
Liu Jicong 已提交
1963 1964
        return NULL;
      }
1965

H
Haojun Liao 已提交
1966
      tscDebug("consumer:0x%" PRIx64 " not ready, retry:%d/40 in 500ms", tmq->consumerId, retryCnt);
L
Liu Jicong 已提交
1967 1968 1969 1970
      taosMsleep(500);
    }
  }

X
Xiaoyu Wang 已提交
1971
  while (1) {
L
Liu Jicong 已提交
1972
    tmqHandleAllDelayedTask(tmq);
1973

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

1978
    rspObj = tmqHandleAllRsp(tmq, timeout, false);
L
Liu Jicong 已提交
1979
    if (rspObj) {
1980
      tscDebug("consumer:0x%" PRIx64 " return rsp %p", tmq->consumerId, rspObj);
L
Liu Jicong 已提交
1981
      return (TAOS_RES*)rspObj;
L
Liu Jicong 已提交
1982
    } else if (terrno == TSDB_CODE_TQ_NO_COMMITTED_OFFSET) {
1983
      tscDebug("consumer:0x%" PRIx64 " return null since no committed offset", tmq->consumerId);
L
Liu Jicong 已提交
1984
      return NULL;
X
Xiaoyu Wang 已提交
1985
    }
1986

1987
    if (timeout >= 0) {
L
Liu Jicong 已提交
1988
      int64_t currentTime = taosGetTimestampMs();
1989 1990 1991
      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 已提交
1992
                 tmq->consumerId, tmq->epoch, startTime, currentTime);
X
Xiaoyu Wang 已提交
1993 1994
        return NULL;
      }
1995
      tsem_timewait(&tmq->rspSem, (timeout - elapsedTime));
L
Liu Jicong 已提交
1996 1997
    } else {
      // use tsem_timewait instead of tsem_wait to avoid unexpected stuck
L
Liu Jicong 已提交
1998
      tsem_timewait(&tmq->rspSem, 1000);
X
Xiaoyu Wang 已提交
1999 2000 2001 2002
    }
  }
}

2003 2004 2005 2006 2007 2008 2009 2010 2011 2012 2013 2014 2015 2016
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);
2017
    }
2018
  }
2019

2020 2021
  tscDebug("consumer:0x%" PRIx64 " rows dist end", pTmq->consumerId);
}
2022

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

2027 2028 2029 2030 2031 2032
  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;
2033 2034 2035
      }
    }

L
Liu Jicong 已提交
2036
    int32_t     retryCnt = 0;
2037
    tmq_list_t* lst = tmq_list_new();
L
Liu Jicong 已提交
2038
    while (1) {
2039
      int32_t rsp = tmq_subscribe(tmq, lst);
L
Liu Jicong 已提交
2040 2041 2042 2043 2044 2045 2046 2047
      if (rsp != TSDB_CODE_MND_CONSUMER_NOT_READY || retryCnt > 5) {
        break;
      } else {
        retryCnt++;
        taosMsleep(500);
      }
    }

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

2053
  taosRemoveRef(tmqMgmt.rsetId, tmq->refId);
L
Liu Jicong 已提交
2054
  return 0;
2055
}
L
Liu Jicong 已提交
2056

L
Liu Jicong 已提交
2057 2058
const char* tmq_err2str(int32_t err) {
  if (err == 0) {
L
Liu Jicong 已提交
2059
    return "success";
L
Liu Jicong 已提交
2060
  } else if (err == -1) {
L
Liu Jicong 已提交
2061 2062 2063
    return "fail";
  } else {
    return tstrerror(err);
L
Liu Jicong 已提交
2064 2065
  }
}
L
Liu Jicong 已提交
2066

L
Liu Jicong 已提交
2067 2068 2069 2070 2071
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;
2072 2073
  } else if (TD_RES_TMQ_METADATA(res)) {
    return TMQ_RES_METADATA;
L
Liu Jicong 已提交
2074 2075 2076 2077 2078
  } else {
    return TMQ_RES_INVALID;
  }
}

L
Liu Jicong 已提交
2079
const char* tmq_get_topic_name(TAOS_RES* res) {
L
Liu Jicong 已提交
2080 2081
  if (TD_RES_TMQ(res)) {
    SMqRspObj* pRspObj = (SMqRspObj*)res;
L
Liu Jicong 已提交
2082
    return strchr(pRspObj->topic, '.') + 1;
L
Liu Jicong 已提交
2083 2084 2085
  } else if (TD_RES_TMQ_META(res)) {
    SMqMetaRspObj* pMetaRspObj = (SMqMetaRspObj*)res;
    return strchr(pMetaRspObj->topic, '.') + 1;
2086 2087 2088
  } else if (TD_RES_TMQ_METADATA(res)) {
    SMqTaosxRspObj* pRspObj = (SMqTaosxRspObj*)res;
    return strchr(pRspObj->topic, '.') + 1;
L
Liu Jicong 已提交
2089 2090 2091 2092 2093
  } else {
    return NULL;
  }
}

L
Liu Jicong 已提交
2094 2095 2096 2097
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 已提交
2098 2099 2100
  } else if (TD_RES_TMQ_META(res)) {
    SMqMetaRspObj* pMetaRspObj = (SMqMetaRspObj*)res;
    return strchr(pMetaRspObj->db, '.') + 1;
2101 2102 2103
  } else if (TD_RES_TMQ_METADATA(res)) {
    SMqTaosxRspObj* pRspObj = (SMqTaosxRspObj*)res;
    return strchr(pRspObj->db, '.') + 1;
L
Liu Jicong 已提交
2104 2105 2106 2107 2108
  } else {
    return NULL;
  }
}

L
Liu Jicong 已提交
2109 2110 2111 2112
int32_t tmq_get_vgroup_id(TAOS_RES* res) {
  if (TD_RES_TMQ(res)) {
    SMqRspObj* pRspObj = (SMqRspObj*)res;
    return pRspObj->vgId;
L
Liu Jicong 已提交
2113 2114 2115
  } else if (TD_RES_TMQ_META(res)) {
    SMqMetaRspObj* pMetaRspObj = (SMqMetaRspObj*)res;
    return pMetaRspObj->vgId;
2116
  } else if (TD_RES_TMQ_METADATA(res)) {
2117 2118
    SMqTaosxRspObj* pRspObj = (SMqTaosxRspObj*)res;
    return pRspObj->vgId;
L
Liu Jicong 已提交
2119 2120 2121 2122
  } else {
    return -1;
  }
}
L
Liu Jicong 已提交
2123

2124 2125 2126 2127 2128 2129 2130 2131 2132 2133 2134 2135 2136 2137 2138 2139 2140 2141 2142 2143 2144 2145 2146
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 已提交
2147 2148 2149 2150 2151 2152 2153
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;
    }
2154
    return (const char*)taosArrayGetP(pRspObj->rsp.blockTbName, pRspObj->resIter);
2155 2156
  } else if (TD_RES_TMQ_METADATA(res)) {
    SMqTaosxRspObj* pRspObj = (SMqTaosxRspObj*)res;
L
Liu Jicong 已提交
2157 2158 2159
    if (!pRspObj->rsp.withTbName || pRspObj->rsp.blockTbName == NULL || pRspObj->resIter < 0 ||
        pRspObj->resIter >= pRspObj->rsp.blockNum) {
      return NULL;
2160
    }
L
Liu Jicong 已提交
2161 2162
    return (const char*)taosArrayGetP(pRspObj->rsp.blockTbName, pRspObj->resIter);
  }
L
Liu Jicong 已提交
2163 2164
  return NULL;
}
2165

2166 2167 2168 2169
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
2170
    asyncCommitOffset(tmq, pRes, TDMT_VND_TMQ_COMMIT_OFFSET, cb, param);
2171
  }
L
Liu Jicong 已提交
2172 2173
}

2174
static void commitCallBackFn(tmq_t *UNUSED_PARAM(tmq), int32_t code, void* param) {
2175 2176 2177
  SSyncCommitInfo* pInfo = (SSyncCommitInfo*) param;
  pInfo->code = code;
  tsem_post(&pInfo->sem);
2178
}
2179

2180 2181 2182 2183 2184 2185 2186 2187 2188
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 已提交
2189
  } else {
2190
    asyncCommitOffset(tmq, pRes, TDMT_VND_TMQ_COMMIT_OFFSET, commitCallBackFn, pInfo);
2191 2192
  }

2193 2194
  tsem_wait(&pInfo->sem);
  code = pInfo->code;
H
Haojun Liao 已提交
2195 2196

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

X
Xiaoyu Wang 已提交
2199
  tscDebug("consumer:0x%" PRIx64 " sync commit done, code:%s", tmq->consumerId, tstrerror(code));
2200 2201 2202 2203 2204 2205 2206 2207 2208 2209 2210 2211
  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);
2212
    doUpdateLocalEp(pTmq, head->epoch, &rsp);
2213 2214 2215
    tDeleteSMqAskEpRsp(&rsp);
  }

H
Haojun Liao 已提交
2216
  tsem_post(&pInfo->sem);
2217 2218 2219 2220 2221 2222 2223 2224 2225 2226 2227 2228 2229 2230 2231 2232 2233 2234 2235 2236 2237 2238 2239 2240 2241 2242
}

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 已提交
2243
  tsem_init(&pInfo->sem, 0, 0);
2244 2245

  asyncAskEp(pTmq, updateEpCallbackFn, pInfo);
H
Haojun Liao 已提交
2246
  tsem_wait(&pInfo->sem);
2247 2248

  int32_t code = pInfo->code;
H
Haojun Liao 已提交
2249
  tsem_destroy(&pInfo->sem);
2250 2251 2252 2253 2254
  taosMemoryFree(pInfo);
  return code;
}

void asyncAskEp(tmq_t* pTmq, __tmq_askep_fn_t askEpFn, void* param) {
2255
  SMqAskEpReq req = {0};
2256 2257 2258
  req.consumerId = pTmq->consumerId;
  req.epoch = pTmq->epoch;
  strcpy(req.cgroup, pTmq->groupId);
2259 2260 2261

  int32_t tlen = tSerializeSMqAskEpReq(NULL, 0, &req);
  if (tlen < 0) {
2262 2263 2264
    tscError("consumer:0x%" PRIx64 ", tSerializeSMqAskEpReq failed", pTmq->consumerId);
    askEpFn(pTmq, TSDB_CODE_INVALID_PARA, NULL, param);
    return;
2265 2266 2267 2268
  }

  void* pReq = taosMemoryCalloc(1, tlen);
  if (pReq == NULL) {
2269 2270 2271
    tscError("consumer:0x%" PRIx64 ", failed to malloc askEpReq msg, size:%d", pTmq->consumerId, tlen);
    askEpFn(pTmq, TSDB_CODE_OUT_OF_MEMORY, NULL, param);
    return;
2272 2273 2274
  }

  if (tSerializeSMqAskEpReq(pReq, tlen, &req) < 0) {
2275
    tscError("consumer:0x%" PRIx64 ", tSerializeSMqAskEpReq %d failed", pTmq->consumerId, tlen);
2276
    taosMemoryFree(pReq);
2277 2278 2279

    askEpFn(pTmq, TSDB_CODE_INVALID_PARA, NULL, param);
    return;
2280 2281 2282 2283
  }

  SMqAskEpCbParam* pParam = taosMemoryCalloc(1, sizeof(SMqAskEpCbParam));
  if (pParam == NULL) {
2284
    tscError("consumer:0x%" PRIx64 ", failed to malloc subscribe param", pTmq->consumerId);
2285
    taosMemoryFree(pReq);
2286 2287 2288

    askEpFn(pTmq, TSDB_CODE_OUT_OF_MEMORY, NULL, param);
    return;
2289 2290
  }

2291 2292 2293 2294
  pParam->refId = pTmq->refId;
  pParam->epoch = pTmq->epoch;
  pParam->pUserFn = askEpFn;
  pParam->pParam = param;
2295 2296 2297 2298 2299

  SMsgSendInfo* sendInfo = taosMemoryCalloc(1, sizeof(SMsgSendInfo));
  if (sendInfo == NULL) {
    taosMemoryFree(pParam);
    taosMemoryFree(pReq);
2300 2301
    askEpFn(pTmq, TSDB_CODE_OUT_OF_MEMORY, NULL, param);
    return;
2302 2303
  }

X
Xiaoyu Wang 已提交
2304
  sendInfo->msgInfo = (SDataBuf){.pData = pReq, .len = tlen, .handle = NULL};
2305 2306 2307 2308

  sendInfo->requestId = generateRequestId();
  sendInfo->requestObjRefId = 0;
  sendInfo->param = pParam;
2309
  sendInfo->fp = askEpCallbackFn;
2310 2311
  sendInfo->msgType = TDMT_MND_TMQ_ASK_EP;

2312 2313
  SEpSet epSet = getEpSet_s(&pTmq->pTscObj->pAppInfo->mgmtEp);
  tscDebug("consumer:0x%" PRIx64 " ask ep from mnode, reqId:0x%" PRIx64, pTmq->consumerId, sendInfo->requestId);
2314 2315

  int64_t transporterId = 0;
2316
  asyncSendMsgToServer(pTmq->pTscObj->pAppInfo->pTransporter, &epSet, &transporterId, sendInfo);
2317 2318 2319 2320 2321 2322 2323
}

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

int32_t tmqCommitDone(SMqCommitCbParamSet* pParamSet) {
2324 2325 2326
  int64_t refId = pParamSet->refId;

  tmq_t* tmq = taosAcquireRef(tmqMgmt.rsetId, refId);
2327 2328 2329 2330 2331 2332 2333
  if (tmq == NULL) {
    taosMemoryFree(pParamSet);
    terrno = TSDB_CODE_TMQ_CONSUMER_CLOSED;
    return -1;
  }

  // if no more waiting rsp
2334
  pParamSet->callbackFn(tmq, pParamSet->code, pParamSet->userParam);
2335
  taosMemoryFree(pParamSet);
2336 2337

  taosReleaseRef(tmqMgmt.rsetId, refId);
2338
  return 0;
2339 2340
}

2341
void commitRspCountDown(SMqCommitCbParamSet* pParamSet, int64_t consumerId, const char* pTopic, int32_t vgId) {
2342 2343
  int32_t waitingRspNum = atomic_sub_fetch_32(&pParamSet->waitingRspNum, 1);
  if (waitingRspNum == 0) {
H
Haojun Liao 已提交
2344 2345
    tscDebug("consumer:0x%" PRIx64 " topic:%s vgId:%d all commit-rsp received, commit completed", consumerId, pTopic,
             vgId);
2346
    tmqCommitDone(pParamSet);
H
Haojun Liao 已提交
2347 2348 2349
  } else {
    tscDebug("consumer:0x%" PRIx64 " topic:%s vgId:%d commit-rsp received, remain:%d", consumerId, pTopic, vgId,
             waitingRspNum);
2350 2351
  }
}
2352 2353 2354 2355 2356 2357 2358 2359 2360 2361 2362 2363 2364 2365 2366 2367 2368 2369 2370 2371 2372 2373

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 已提交
2374 2375
}

2376 2377 2378 2379 2380 2381 2382 2383 2384 2385 2386 2387 2388 2389 2390 2391 2392 2393 2394 2395 2396
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,
2397
                                       .vgId = pParam->vgId};
2398 2399 2400 2401 2402 2403 2404 2405 2406 2407 2408 2409 2410 2411 2412 2413 2414 2415 2416 2417 2418 2419

    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 已提交
2420
int32_t tmq_get_topic_assignment(tmq_t* tmq, const char* pTopicName, tmq_topic_assignment** assignment,
H
Haojun Liao 已提交
2421
                                 int32_t* numOfAssignment) {
H
Haojun Liao 已提交
2422 2423 2424
  *numOfAssignment = 0;
  *assignment = NULL;

2425
  int32_t accId = tmq->pTscObj->acctId;
2426
  char    tname[128] = {0};
2427 2428 2429
  sprintf(tname, "%d.%s", accId, pTopicName);

  SMqClientTopic* pTopic = getTopicByName(tmq, tname);
H
Haojun Liao 已提交
2430 2431 2432 2433 2434 2435
  if (pTopic == NULL) {
    return TSDB_CODE_INVALID_PARA;
  }

  // in case of snapshot is opened, no valid offset will return
  *numOfAssignment = taosArrayGetSize(pTopic->vgs);
2436 2437 2438 2439 2440 2441 2442 2443

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

2444 2445
  bool needFetch = false;

H
Haojun Liao 已提交
2446 2447
  for (int32_t j = 0; j < (*numOfAssignment); ++j) {
    SMqClientVg* pClientVg = taosArrayGet(pTopic->vgs, j);
2448
    if (!pClientVg->receivedInfoFromVnode) {
2449 2450 2451
      needFetch = true;
      break;
    }
H
Haojun Liao 已提交
2452 2453 2454 2455 2456 2457 2458 2459 2460 2461

    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;
2462
    pAssignment->vgId = pClientVg->vgId;
H
Haojun Liao 已提交
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 2535 2536 2537 2538 2539 2540 2541 2542 2543 2544 2545 2546
  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);
2547
      *assignment = NULL;
2548 2549 2550 2551 2552 2553 2554 2555 2556
      *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;
    }

2557 2558 2559 2560 2561 2562 2563 2564 2565 2566 2567 2568 2569 2570 2571 2572 2573 2574 2575 2576 2577 2578 2579 2580 2581
    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;
      }
    }

2582 2583 2584 2585 2586
    destroyCommonInfo(pCommon);
    return code;
  } else {
    return TSDB_CODE_SUCCESS;
  }
H
Haojun Liao 已提交
2587 2588
}

T
t_max 已提交
2589 2590 2591 2592 2593 2594 2595 2596
void tmq_free_assignment(tmq_topic_assignment* pAssignment) {
    if (pAssignment == NULL) {
        return;
    }

    taosMemoryFree(pAssignment);
}

2597
int32_t tmq_offset_seek(tmq_t* tmq, const char* pTopicName, int32_t vgId, int64_t offset) {
H
Haojun Liao 已提交
2598
  if (tmq == NULL) {
H
Haojun Liao 已提交
2599
    tscError("invalid tmq handle, null");
H
Haojun Liao 已提交
2600 2601 2602
    return TSDB_CODE_INVALID_PARA;
  }

2603 2604 2605 2606 2607
  int32_t accId = tmq->pTscObj->acctId;
  char tname[128] = {0};
  sprintf(tname, "%d.%s", accId, pTopicName);

  SMqClientTopic* pTopic = getTopicByName(tmq, tname);
H
Haojun Liao 已提交
2608
  if (pTopic == NULL) {
2609
    tscError("consumer:0x%" PRIx64 " invalid topic name:%s", tmq->consumerId, pTopicName);
H
Haojun Liao 已提交
2610 2611 2612 2613
    return TSDB_CODE_INVALID_PARA;
  }

  SMqClientVg* pVg = NULL;
H
Haojun Liao 已提交
2614 2615
  int32_t      numOfVgs = taosArrayGetSize(pTopic->vgs);
  for (int32_t i = 0; i < numOfVgs; ++i) {
H
Haojun Liao 已提交
2616
    SMqClientVg* pClientVg = taosArrayGet(pTopic->vgs, i);
2617
    if (pClientVg->vgId == vgId) {
H
Haojun Liao 已提交
2618 2619 2620 2621 2622 2623
      pVg = pClientVg;
      break;
    }
  }

  if (pVg == NULL) {
2624
    tscError("consumer:0x%" PRIx64 " invalid vgroup id:%d", tmq->consumerId, vgId);
H
Haojun Liao 已提交
2625 2626 2627 2628 2629 2630
    return TSDB_CODE_INVALID_PARA;
  }

  SVgOffsetInfo* pOffsetInfo = &pVg->offsetInfo;

  int32_t type = pOffsetInfo->currentOffset.type;
2631
  if (type != TMQ_OFFSET__LOG && !OFFSET_IS_RESET_OFFSET(type)) {
2632
    tscError("consumer:0x%" PRIx64 " offset type:%d not wal version, seek not allowed", tmq->consumerId, type);
H
Haojun Liao 已提交
2633 2634 2635
    return TSDB_CODE_INVALID_PARA;
  }

2636
  if (type == TMQ_OFFSET__LOG && (offset < pOffsetInfo->walVerBegin || offset > pOffsetInfo->walVerEnd)) {
H
Haojun Liao 已提交
2637 2638
    tscError("consumer:0x%" PRIx64 " invalid seek params, offset:%" PRId64 ", valid range:[%" PRId64 ", %" PRId64 "]",
             tmq->consumerId, offset, pOffsetInfo->walVerBegin, pOffsetInfo->walVerEnd);
H
Haojun Liao 已提交
2639 2640 2641
    return TSDB_CODE_INVALID_PARA;
  }

H
Haojun Liao 已提交
2642 2643 2644
  // update the offset, and then commit to vnode
  if (pOffsetInfo->currentOffset.type == TMQ_OFFSET__LOG) {
    pOffsetInfo->currentOffset.version = offset;
2645
    pOffsetInfo->committedOffset.version = INT64_MIN;
2646
    pVg->seekUpdated = true;
H
Haojun Liao 已提交
2647 2648
  }

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

H
Haojun Liao 已提交
2652
  tscDebug("consumer:0x%" PRIx64 " seek to %" PRId64 " on vgId:%d", tmq->consumerId, offset, pVg->vgId);
2653 2654 2655 2656 2657 2658 2659 2660 2661 2662 2663 2664 2665 2666 2667 2668 2669 2670

  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 已提交
2671 2672 2673 2674 2675 2676
  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;
2677
}