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

L
Liu Jicong 已提交
16
#include "cJSON.h"
17 18 19
#include "clientInt.h"
#include "clientLog.h"
#include "parser.h"
H
Haojun Liao 已提交
20
#include "tdatablock.h"
21 22
#include "tdef.h"
#include "tglobal.h"
X
Xiaoyu Wang 已提交
23
#include "tqueue.h"
24
#include "tref.h"
L
Liu Jicong 已提交
25 26
#include "ttimer.h"

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

L
Liu Jicong 已提交
136
typedef struct {
H
Haojun Liao 已提交
137
  int64_t      pollCnt;
138
  int64_t      numOfRows;
L
Liu Jicong 已提交
139 140
  STqOffsetVal committedOffset;
  STqOffsetVal currentOffset;
H
Haojun Liao 已提交
141 142 143
  int32_t      vgId;
  int32_t      vgStatus;
  int32_t      vgSkipCnt;
X
Xiaoyu Wang 已提交
144
  int64_t      emptyBlockReceiveTs;  // once empty block is received, idle for ignoreCnt then start to poll data
H
Haojun Liao 已提交
145
  SEpSet       epSet;
146 147
} SMqClientVg;

L
Liu Jicong 已提交
148
typedef struct {
149 150 151
  char           topicName[TSDB_TOPIC_FNAME_LEN];
  char           db[TSDB_DB_FNAME_LEN];
  SArray*        vgs;  // SArray<SMqClientVg>
L
Liu Jicong 已提交
152
  SSchemaWrapper schema;
153 154
} SMqClientTopic;

L
Liu Jicong 已提交
155 156
typedef struct {
  int8_t          tmqRspType;
157
  int32_t         epoch;  // epoch can be used to guard the vgHandle
158
  int32_t         vgId;
L
Liu Jicong 已提交
159 160
  SMqClientVg*    vgHandle;
  SMqClientTopic* topicHandle;
H
Haojun Liao 已提交
161
  uint64_t        reqId;
162
  SEpSet*         pEpset;
L
Liu Jicong 已提交
163
  union {
L
Liu Jicong 已提交
164 165
    SMqDataRsp dataRsp;
    SMqMetaRsp metaRsp;
L
Liu Jicong 已提交
166
    STaosxRsp  taosxRsp;
L
Liu Jicong 已提交
167
  };
L
Liu Jicong 已提交
168 169
} SMqPollRspWrapper;

L
Liu Jicong 已提交
170
typedef struct {
171 172
  int64_t refId;
  int32_t epoch;
L
Liu Jicong 已提交
173 174
  tsem_t  rspSem;
  int32_t rspErr;
L
Liu Jicong 已提交
175
} SMqSubscribeCbParam;
L
Liu Jicong 已提交
176

L
Liu Jicong 已提交
177
typedef struct {
178 179 180 181
  int64_t          refId;
  int32_t          epoch;
  void*            pParam;
  __tmq_askep_fn_t pUserFn;
182 183
} SMqAskEpCbParam;

L
Liu Jicong 已提交
184
typedef struct {
185 186
  int64_t         refId;
  int32_t         epoch;
L
Liu Jicong 已提交
187
  SMqClientVg*    pVg;
L
Liu Jicong 已提交
188
  SMqClientTopic* pTopic;
L
Liu Jicong 已提交
189
  int32_t         vgId;
X
Xiaoyu Wang 已提交
190
  uint64_t        requestId;  // request id for debug purpose
X
Xiaoyu Wang 已提交
191
} SMqPollCbParam;
192

193
typedef struct {
194 195
  int64_t        refId;
  int32_t        epoch;
L
Liu Jicong 已提交
196 197
  int32_t        waitingRspNum;
  int32_t        totalRspNum;
198
  int32_t        code;
199
  tmq_commit_cb* callbackFn;
L
Liu Jicong 已提交
200 201
  /*SArray*        successfulOffsets;*/
  /*SArray*        failedOffsets;*/
X
Xiaoyu Wang 已提交
202
  void* userParam;
203 204 205 206 207
} SMqCommitCbParamSet;

typedef struct {
  SMqCommitCbParamSet* params;
  STqOffset*           pOffset;
H
Haojun Liao 已提交
208 209 210
  char                 topicName[TSDB_TOPIC_FNAME_LEN];
  int32_t              vgId;
  tmq_t*               pTmq;
211
} SMqCommitCbParam;
212

213 214 215 216 217
typedef struct SSyncCommitInfo {
  tsem_t  sem;
  int32_t code;
} SSyncCommitInfo;

218
static int32_t doAskEp(tmq_t* tmq);
219 220
static int32_t makeTopicVgroupKey(char* dst, const char* topicName, int32_t vg);
static int32_t tmqCommitDone(SMqCommitCbParamSet* pParamSet);
221 222
static int32_t doSendCommitMsg(tmq_t* tmq, SMqClientVg* pVg, const char* pTopicName, SMqCommitCbParamSet* pParamSet,
                               int32_t index, int32_t totalVgroups);
X
Xiaoyu Wang 已提交
223 224 225
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);
226

227
tmq_conf_t* tmq_conf_new() {
wafwerar's avatar
wafwerar 已提交
228
  tmq_conf_t* conf = taosMemoryCalloc(1, sizeof(tmq_conf_t));
229 230 231 232 233
  if (conf == NULL) {
    terrno = TSDB_CODE_OUT_OF_MEMORY;
    return conf;
  }

234
  conf->withTbName = false;
L
Liu Jicong 已提交
235
  conf->autoCommit = true;
236
  conf->autoCommitInterval = DEFAULT_AUTO_COMMIT_INTERVAL;
237
  conf->resetOffset = TMQ_OFFSET__RESET_EARLIEAST;
238
  conf->hbBgEnable = true;
239

240 241 242
  return conf;
}

L
Liu Jicong 已提交
243
void tmq_conf_destroy(tmq_conf_t* conf) {
L
Liu Jicong 已提交
244
  if (conf) {
245 246 247 248 249 250 251 252 253
    if (conf->ip) {
      taosMemoryFree(conf->ip);
    }
    if (conf->user) {
      taosMemoryFree(conf->user);
    }
    if (conf->pass) {
      taosMemoryFree(conf->pass);
    }
L
Liu Jicong 已提交
254 255
    taosMemoryFree(conf);
  }
L
Liu Jicong 已提交
256 257 258
}

tmq_conf_res_t tmq_conf_set(tmq_conf_t* conf, const char* key, const char* value) {
259
  if (strcasecmp(key, "group.id") == 0) {
L
Liu Jicong 已提交
260
    tstrncpy(conf->groupId, value, TSDB_CGROUP_LEN);
L
Liu Jicong 已提交
261
    return TMQ_CONF_OK;
262
  }
L
Liu Jicong 已提交
263

264
  if (strcasecmp(key, "client.id") == 0) {
L
Liu Jicong 已提交
265
    tstrncpy(conf->clientId, value, 256);
L
Liu Jicong 已提交
266 267
    return TMQ_CONF_OK;
  }
L
Liu Jicong 已提交
268

269 270
  if (strcasecmp(key, "enable.auto.commit") == 0) {
    if (strcasecmp(value, "true") == 0) {
L
Liu Jicong 已提交
271
      conf->autoCommit = true;
L
Liu Jicong 已提交
272
      return TMQ_CONF_OK;
273
    } else if (strcasecmp(value, "false") == 0) {
L
Liu Jicong 已提交
274
      conf->autoCommit = false;
L
Liu Jicong 已提交
275 276 277 278
      return TMQ_CONF_OK;
    } else {
      return TMQ_CONF_INVALID;
    }
279
  }
L
Liu Jicong 已提交
280

281
  if (strcasecmp(key, "auto.commit.interval.ms") == 0) {
282
    conf->autoCommitInterval = taosStr2int64(value);
L
Liu Jicong 已提交
283 284 285
    return TMQ_CONF_OK;
  }

286 287 288
  if (strcasecmp(key, "auto.offset.reset") == 0) {
    if (strcasecmp(value, "none") == 0) {
      conf->resetOffset = TMQ_OFFSET__RESET_NONE;
L
Liu Jicong 已提交
289
      return TMQ_CONF_OK;
290 291
    } else if (strcasecmp(value, "earliest") == 0) {
      conf->resetOffset = TMQ_OFFSET__RESET_EARLIEAST;
L
Liu Jicong 已提交
292
      return TMQ_CONF_OK;
293 294
    } else if (strcasecmp(value, "latest") == 0) {
      conf->resetOffset = TMQ_OFFSET__RESET_LATEST;
L
Liu Jicong 已提交
295 296 297 298 299
      return TMQ_CONF_OK;
    } else {
      return TMQ_CONF_INVALID;
    }
  }
L
Liu Jicong 已提交
300

301 302
  if (strcasecmp(key, "msg.with.table.name") == 0) {
    if (strcasecmp(value, "true") == 0) {
303
      conf->withTbName = true;
L
Liu Jicong 已提交
304
      return TMQ_CONF_OK;
305
    } else if (strcasecmp(value, "false") == 0) {
306
      conf->withTbName = false;
L
Liu Jicong 已提交
307
      return TMQ_CONF_OK;
308 309 310 311 312
    } else {
      return TMQ_CONF_INVALID;
    }
  }

313 314
  if (strcasecmp(key, "experimental.snapshot.enable") == 0) {
    if (strcasecmp(value, "true") == 0) {
L
Liu Jicong 已提交
315
      conf->snapEnable = true;
316
      return TMQ_CONF_OK;
317
    } else if (strcasecmp(value, "false") == 0) {
L
Liu Jicong 已提交
318
      conf->snapEnable = false;
319 320 321 322 323 324
      return TMQ_CONF_OK;
    } else {
      return TMQ_CONF_INVALID;
    }
  }

325
  if (strcasecmp(key, "experimental.snapshot.batch.size") == 0) {
326
    conf->snapBatchSize = taosStr2int64(value);
L
Liu Jicong 已提交
327 328 329
    return TMQ_CONF_OK;
  }

330
  if (strcasecmp(key, "enable.heartbeat.background") == 0) {
X
Xiaoyu Wang 已提交
331 332 333 334 335 336 337 338 339 340
    //    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 已提交
341 342
  }

343
  if (strcasecmp(key, "td.connect.ip") == 0) {
344
    conf->ip = taosStrdup(value);
L
Liu Jicong 已提交
345 346
    return TMQ_CONF_OK;
  }
347

348
  if (strcasecmp(key, "td.connect.user") == 0) {
349
    conf->user = taosStrdup(value);
L
Liu Jicong 已提交
350 351
    return TMQ_CONF_OK;
  }
352

353
  if (strcasecmp(key, "td.connect.pass") == 0) {
354
    conf->pass = taosStrdup(value);
L
Liu Jicong 已提交
355 356
    return TMQ_CONF_OK;
  }
357

358
  if (strcasecmp(key, "td.connect.port") == 0) {
359
    conf->port = taosStr2int64(value);
L
Liu Jicong 已提交
360 361
    return TMQ_CONF_OK;
  }
362

363
  if (strcasecmp(key, "td.connect.db") == 0) {
L
Liu Jicong 已提交
364 365 366
    return TMQ_CONF_OK;
  }

L
Liu Jicong 已提交
367
  return TMQ_CONF_UNKNOWN;
368 369
}

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

L
Liu Jicong 已提交
372 373
int32_t tmq_list_append(tmq_list_t* list, const char* src) {
  SArray* container = &list->container;
374
  if (src == NULL || src[0] == 0) return -1;
375
  char* topic = taosStrdup(src);
L
fix  
Liu Jicong 已提交
376
  if (taosArrayPush(container, &topic) == NULL) return -1;
377 378 379
  return 0;
}

L
Liu Jicong 已提交
380
void tmq_list_destroy(tmq_list_t* list) {
L
Liu Jicong 已提交
381
  SArray* container = &list->container;
L
Liu Jicong 已提交
382
  taosArrayDestroyP(container, taosMemoryFree);
L
Liu Jicong 已提交
383 384
}

L
Liu Jicong 已提交
385 386 387 388 389 390 391 392 393 394
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 已提交
395 396
static SMqClientVg* foundClientVg(SArray* pTopicList, const char* pName, int32_t vgId, int32_t* index,
                                  int32_t* numOfVgroups) {
397 398 399 400
  int32_t numOfTopics = taosArrayGetSize(pTopicList);
  *index = -1;
  *numOfVgroups = 0;

X
Xiaoyu Wang 已提交
401
  for (int32_t i = 0; i < numOfTopics; ++i) {
402 403
    SMqClientTopic* pTopic = taosArrayGet(pTopicList, i);
    if (strcmp(pTopic->topicName, pName) != 0) {
404 405 406
      continue;
    }

407 408
    *numOfVgroups = taosArrayGetSize(pTopic->vgs);
    for (int32_t j = 0; j < (*numOfVgroups); ++j) {
409
      SMqClientVg* pClientVg = taosArrayGet(pTopic->vgs, j);
410 411 412
      if (pClientVg->vgId == vgId) {
        *index = j;
        return pClientVg;
413 414
      }
    }
L
Liu Jicong 已提交
415
  }
416 417

  return NULL;
L
Liu Jicong 已提交
418
}
419

420 421 422
// 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 已提交
423
static int32_t tmqCommitCb(void* param, SDataBuf* pBuf, int32_t code) {
424
  SMqCommitCbParam*    pParam = (SMqCommitCbParam*)param;
425
  SMqCommitCbParamSet* pParamSet = (SMqCommitCbParamSet*)pParam->params;
H
Haojun Liao 已提交
426

X
Xiaoyu Wang 已提交
427 428 429 430 431 432 433 434 435 436 437 438 439 440 441 442 443 444 445 446 447 448 449 450 451 452 453 454 455 456
  //  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",
  //                 pParam->pTmq->consumerId, pParam->topicName, pVg->vgId, pVg->committedOffset.version,
  //                 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
457

L
Liu Jicong 已提交
458
  taosMemoryFree(pParam->pOffset);
L
Liu Jicong 已提交
459
  taosMemoryFree(pBuf->pData);
dengyihao's avatar
dengyihao 已提交
460
  taosMemoryFree(pBuf->pEpSet);
L
Liu Jicong 已提交
461

462
  commitRspCountDown(pParamSet, pParam->pTmq->consumerId, pParam->topicName, pParam->vgId);
463 464 465
  return 0;
}

466 467
static int32_t doSendCommitMsg(tmq_t* tmq, SMqClientVg* pVg, const char* pTopicName, SMqCommitCbParamSet* pParamSet,
                               int32_t index, int32_t totalVgroups) {
L
Liu Jicong 已提交
468 469
  STqOffset* pOffset = taosMemoryCalloc(1, sizeof(STqOffset));
  if (pOffset == NULL) {
470
    return TSDB_CODE_OUT_OF_MEMORY;
L
Liu Jicong 已提交
471
  }
472

L
Liu Jicong 已提交
473
  pOffset->val = pVg->currentOffset;
474

L
Liu Jicong 已提交
475 476 477
  int32_t groupLen = strlen(tmq->groupId);
  memcpy(pOffset->subKey, tmq->groupId, groupLen);
  pOffset->subKey[groupLen] = TMQ_SEPARATOR;
H
Haojun Liao 已提交
478
  strcpy(pOffset->subKey + groupLen + 1, pTopicName);
L
Liu Jicong 已提交
479

480 481
  int32_t len = 0;
  int32_t code = 0;
L
Liu Jicong 已提交
482 483
  tEncodeSize(tEncodeSTqOffset, pOffset, len, code);
  if (code < 0) {
484
    return TSDB_CODE_INVALID_PARA;
L
Liu Jicong 已提交
485
  }
486

L
Liu Jicong 已提交
487
  void* buf = taosMemoryCalloc(1, sizeof(SMsgHead) + len);
L
Liu Jicong 已提交
488 489
  if (buf == NULL) {
    taosMemoryFree(pOffset);
490
    return TSDB_CODE_OUT_OF_MEMORY;
L
Liu Jicong 已提交
491
  }
492

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

L
Liu Jicong 已提交
495 496 497 498 499
  void* abuf = POINTER_SHIFT(buf, sizeof(SMsgHead));

  SEncoder encoder;
  tEncoderInit(&encoder, abuf, len);
  tEncodeSTqOffset(&encoder, pOffset);
L
Liu Jicong 已提交
500
  tEncoderClear(&encoder);
L
Liu Jicong 已提交
501 502

  // build param
503
  SMqCommitCbParam* pParam = taosMemoryCalloc(1, sizeof(SMqCommitCbParam));
L
Liu Jicong 已提交
504
  if (pParam == NULL) {
L
Liu Jicong 已提交
505
    taosMemoryFree(pOffset);
L
Liu Jicong 已提交
506
    taosMemoryFree(buf);
507
    return TSDB_CODE_OUT_OF_MEMORY;
L
Liu Jicong 已提交
508
  }
509

L
Liu Jicong 已提交
510 511
  pParam->params = pParamSet;
  pParam->pOffset = pOffset;
H
Haojun Liao 已提交
512 513 514
  pParam->vgId = pVg->vgId;
  pParam->pTmq = tmq;

H
Haojun Liao 已提交
515
  tstrncpy(pParam->topicName, pTopicName, tListLen(pParam->topicName));
L
Liu Jicong 已提交
516 517 518 519

  // build send info
  SMsgSendInfo* pMsgSendInfo = taosMemoryCalloc(1, sizeof(SMsgSendInfo));
  if (pMsgSendInfo == NULL) {
L
Liu Jicong 已提交
520
    taosMemoryFree(pOffset);
L
Liu Jicong 已提交
521 522
    taosMemoryFree(buf);
    taosMemoryFree(pParam);
523
    return TSDB_CODE_OUT_OF_MEMORY;
L
Liu Jicong 已提交
524
  }
525

526
  pMsgSendInfo->msgInfo = (SDataBuf) { .pData = buf, .len = sizeof(SMsgHead) + len, .handle = NULL };
L
Liu Jicong 已提交
527 528 529 530

  pMsgSendInfo->requestId = generateRequestId();
  pMsgSendInfo->requestObjRefId = 0;
  pMsgSendInfo->param = pParam;
L
Liu Jicong 已提交
531
  pMsgSendInfo->paramFreeFp = taosMemoryFree;
532
  pMsgSendInfo->fp = tmqCommitCb;
L
Liu Jicong 已提交
533
  pMsgSendInfo->msgType = TDMT_VND_TMQ_COMMIT_OFFSET;
L
Liu Jicong 已提交
534

L
Liu Jicong 已提交
535 536 537
  atomic_add_fetch_32(&pParamSet->waitingRspNum, 1);
  atomic_add_fetch_32(&pParamSet->totalRspNum, 1);

H
Haojun Liao 已提交
538
  SEp* pEp = GET_ACTIVE_EP(&pVg->epSet);
539 540 541 542 543 544 545 546
  char offsetBuf[80] = {0};
  tFormatOffset(offsetBuf, tListLen(offsetBuf), &pOffset->val);

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

L
Liu Jicong 已提交
548 549
  int64_t transporterId = 0;
  asyncSendMsgToServer(tmq->pTscObj->pAppInfo->pTransporter, &pVg->epSet, &transporterId, pMsgSendInfo);
550 551

  return TSDB_CODE_SUCCESS;
L
Liu Jicong 已提交
552 553
}

554 555 556 557 558 559 560 561 562 563 564 565 566
static void asyncCommitOffset(tmq_t* tmq, const TAOS_RES* pRes, tmq_commit_cb* pCommitFp, void* userParam) {
  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 已提交
567
    vgId = pRspObj->vgId;
568 569 570
  } else if (TD_RES_TMQ_META(pRes)) {
    SMqMetaRspObj* pMetaRspObj = (SMqMetaRspObj*)pRes;
    pTopicName = pMetaRspObj->topic;
L
Liu Jicong 已提交
571
    vgId = pMetaRspObj->vgId;
572 573 574
  } else if (TD_RES_TMQ_METADATA(pRes)) {
    SMqTaosxRspObj* pRspObj = (SMqTaosxRspObj*)pRes;
    pTopicName = pRspObj->topic;
575
    vgId = pRspObj->vgId;
L
Liu Jicong 已提交
576
  } else {
577 578
    pCommitFp(tmq, TSDB_CODE_TMQ_INVALID_MSG, userParam);
    return;
L
Liu Jicong 已提交
579 580 581 582
  }

  SMqCommitCbParamSet* pParamSet = taosMemoryCalloc(1, sizeof(SMqCommitCbParamSet));
  if (pParamSet == NULL) {
583 584
    pCommitFp(tmq, TSDB_CODE_OUT_OF_MEMORY, userParam);
    return;
L
Liu Jicong 已提交
585
  }
H
Haojun Liao 已提交
586

587 588
  pParamSet->refId = tmq->refId;
  pParamSet->epoch = tmq->epoch;
589
  pParamSet->callbackFn = pCommitFp;
L
Liu Jicong 已提交
590
  pParamSet->userParam = userParam;
L
Liu Jicong 已提交
591

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

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

  int32_t i = 0;
  for (; i < numOfTopics; i++) {
L
Liu Jicong 已提交
598
    SMqClientTopic* pTopic = taosArrayGet(tmq->clientTopics, i);
599 600
    if (strcmp(pTopic->topicName, pTopicName) == 0) {
      break;
601
    }
602
  }
603

604
  if (i == numOfTopics) {
X
Xiaoyu Wang 已提交
605 606
    tscWarn("consumer:0x%" PRIx64 " failed to find the specified topic:%s, total topics:%d", tmq->consumerId,
            pTopicName, numOfTopics);
607 608 609 610
    taosMemoryFree(pParamSet);
    pCommitFp(tmq, TSDB_CODE_SUCCESS, userParam);
    return;
  }
L
Liu Jicong 已提交
611

612 613 614 615 616 617 618 619
  SMqClientTopic* pTopic = taosArrayGet(tmq->clientTopics, i);

  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 已提交
620
    }
L
Liu Jicong 已提交
621
  }
L
Liu Jicong 已提交
622

623
  if (j == numOfVgroups) {
X
Xiaoyu Wang 已提交
624 625
    tscWarn("consumer:0x%" PRIx64 " failed to find the specified vgId:%d, total Vgs:%d, topic:%s", tmq->consumerId,
            vgId, numOfVgroups, pTopicName);
L
Liu Jicong 已提交
626
    taosMemoryFree(pParamSet);
627 628
    pCommitFp(tmq, TSDB_CODE_SUCCESS, userParam);
    return;
L
Liu Jicong 已提交
629 630
  }

631 632 633
  SMqClientVg* pVg = taosArrayGet(pTopic->vgs, j);
  if (pVg->currentOffset.type > 0 && !tOffsetEqual(&pVg->currentOffset, &pVg->committedOffset)) {
    code = doSendCommitMsg(tmq, pVg, pTopic->topicName, pParamSet, j, numOfVgroups);
L
Liu Jicong 已提交
634

635 636 637 638 639
    // failed to commit, callback user function directly.
    if (code != TSDB_CODE_SUCCESS) {
      taosMemoryFree(pParamSet);
      pCommitFp(tmq, code, userParam);
    }
X
Xiaoyu Wang 已提交
640
  } else {  // do not perform commit, callback user function directly.
641 642
    taosMemoryFree(pParamSet);
    pCommitFp(tmq, code, userParam);
L
Liu Jicong 已提交
643 644 645
  }
}

646
static void asyncCommitAllOffsets(tmq_t* tmq, tmq_commit_cb* pCommitFp, void* userParam) {
647 648
  SMqCommitCbParamSet* pParamSet = taosMemoryCalloc(1, sizeof(SMqCommitCbParamSet));
  if (pParamSet == NULL) {
649 650
    pCommitFp(tmq, TSDB_CODE_OUT_OF_MEMORY, userParam);
    return;
651
  }
652 653 654

  pParamSet->refId = tmq->refId;
  pParamSet->epoch = tmq->epoch;
655
  pParamSet->callbackFn = pCommitFp;
656 657
  pParamSet->userParam = userParam;

658 659 660
  // init as 1 to prevent concurrency issue
  pParamSet->waitingRspNum = 1;

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

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

668 669
    tscDebug("consumer:0x%" PRIx64 " commit offset for topics:%s, numOfVgs:%d", tmq->consumerId, pTopic->topicName,
             numOfVgroups);
670
    for (int32_t j = 0; j < numOfVgroups; j++) {
671 672 673
      SMqClientVg* pVg = taosArrayGet(pTopic->vgs, j);

      if (pVg->currentOffset.type > 0 && !tOffsetEqual(&pVg->currentOffset, &pVg->committedOffset)) {
674
        int32_t code = doSendCommitMsg(tmq, pVg, pTopic->topicName, pParamSet, j, numOfVgroups);
675 676 677 678
        if (code != TSDB_CODE_SUCCESS) {
          tscError("consumer:0x%" PRIx64 " topic:%s vgId:%d offset:%" PRId64 " failed, code:%s ordinal:%d/%d",
                   tmq->consumerId, pTopic->topicName, pVg->vgId, pVg->committedOffset.version, tstrerror(terrno),
                   j + 1, numOfVgroups);
L
Liu Jicong 已提交
679 680
          continue;
        }
H
Haojun Liao 已提交
681 682 683

        // update the offset value.
        pVg->committedOffset = pVg->currentOffset;
684
      } else {
D
dapan1121 已提交
685
        tscDebug("consumer:0x%" PRIx64 " topic:%s vgId:%d, no commit, current:%" PRId64 ", ordinal:%d/%d",
686
                 tmq->consumerId, pTopic->topicName, pVg->vgId, pVg->currentOffset.version, j + 1, numOfVgroups);
687 688 689 690
      }
    }
  }

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

L
Liu Jicong 已提交
694
  // no request is sent
L
Liu Jicong 已提交
695 696
  if (pParamSet->totalRspNum == 0) {
    taosMemoryFree(pParamSet);
697 698
    pCommitFp(tmq, TSDB_CODE_SUCCESS, userParam);
    return;
L
Liu Jicong 已提交
699 700
  }

L
Liu Jicong 已提交
701
  // count down since waiting rsp num init as 1
702
  commitRspCountDown(pParamSet, tmq->consumerId, "", 0);
703 704
}

705 706
static void generateTimedTask(int64_t refId, int32_t type) {
  tmq_t* tmq = taosAcquireRef(tmqMgmt.rsetId, refId);
707
  if (tmq != NULL) {
S
Shengliang Guan 已提交
708
    int8_t* pTaskType = taosAllocateQitem(sizeof(int8_t), DEF_QITEM, 0);
709
    *pTaskType = type;
710 711 712
    taosWriteQitem(tmq->delayedTask, pTaskType);
    tsem_post(&tmq->rspSem);
  }
713
  taosReleaseRef(tmqMgmt.rsetId, refId);
714 715 716 717 718
}

void tmqAssignAskEpTask(void* param, void* tmrId) {
  int64_t refId = *(int64_t*)param;
  generateTimedTask(refId, TMQ_DELAYED_TASK__ASK_EP);
719
  taosMemoryFree(param);
L
Liu Jicong 已提交
720 721 722
}

void tmqAssignDelayedCommitTask(void* param, void* tmrId) {
723
  int64_t refId = *(int64_t*)param;
724
  generateTimedTask(refId, TMQ_DELAYED_TASK__COMMIT);
725
  taosMemoryFree(param);
L
Liu Jicong 已提交
726 727 728
}

void tmqAssignDelayedReportTask(void* param, void* tmrId) {
729 730 731
  int64_t refId = *(int64_t*)param;
  tmq_t*  tmq = taosAcquireRef(tmqMgmt.rsetId, refId);
  if (tmq != NULL) {
S
Shengliang Guan 已提交
732
    int8_t* pTaskType = taosAllocateQitem(sizeof(int8_t), DEF_QITEM, 0);
733 734 735 736
    *pTaskType = TMQ_DELAYED_TASK__REPORT;
    taosWriteQitem(tmq->delayedTask, pTaskType);
    tsem_post(&tmq->rspSem);
  }
737 738

  taosReleaseRef(tmqMgmt.rsetId, refId);
739
  taosMemoryFree(param);
L
Liu Jicong 已提交
740 741
}

742
int32_t tmqHbCb(void* param, SDataBuf* pMsg, int32_t code) {
dengyihao's avatar
dengyihao 已提交
743 744 745 746
  if (pMsg) {
    taosMemoryFree(pMsg->pData);
    taosMemoryFree(pMsg->pEpSet);
  }
747 748 749 750
  return 0;
}

void tmqSendHbReq(void* param, void* tmrId) {
751
  int64_t refId = *(int64_t*)param;
752

X
Xiaoyu Wang 已提交
753
  tmq_t* tmq = taosAcquireRef(tmqMgmt.rsetId, refId);
754
  if (tmq == NULL) {
L
Liu Jicong 已提交
755
    taosMemoryFree(param);
756 757
    return;
  }
D
dapan1121 已提交
758 759 760 761 762

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

L
Liu Jicong 已提交
763
  int32_t tlen = tSerializeSMqHbReq(NULL, 0, &req);
D
dapan1121 已提交
764 765
  if (tlen < 0) {
    tscError("tSerializeSMqHbReq failed");
766
    goto OVER;
D
dapan1121 已提交
767
  }
768

L
Liu Jicong 已提交
769
  void* pReq = taosMemoryCalloc(1, tlen);
D
dapan1121 已提交
770 771
  if (tlen < 0) {
    tscError("failed to malloc MqHbReq msg, size:%d", tlen);
772
    goto OVER;
D
dapan1121 已提交
773
  }
774

D
dapan1121 已提交
775 776 777
  if (tSerializeSMqHbReq(pReq, tlen, &req) < 0) {
    tscError("tSerializeSMqHbReq %d failed", tlen);
    taosMemoryFree(pReq);
778
    goto OVER;
D
dapan1121 已提交
779
  }
780 781 782 783

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

787
  sendInfo->msgInfo = (SDataBuf){ .pData = pReq, .len = tlen, .handle = NULL };
788 789 790 791 792

  sendInfo->requestId = generateRequestId();
  sendInfo->requestObjRefId = 0;
  sendInfo->param = NULL;
  sendInfo->fp = tmqHbCb;
L
Liu Jicong 已提交
793
  sendInfo->msgType = TDMT_MND_TMQ_HB;
794 795 796 797 798 799 800

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

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

OVER:
801
  taosTmrReset(tmqSendHbReq, 1000, param, tmqMgmt.timer, &tmq->hbLiveTimer);
802
  taosReleaseRef(tmqMgmt.rsetId, refId);
803 804
}

805 806
static void defaultCommitCbFn(tmq_t* pTmq, int32_t code, void* param) {
  if (code != 0) {
X
Xiaoyu Wang 已提交
807
    tscDebug("consumer:0x%" PRIx64 ", failed to commit offset, code:%s", pTmq->consumerId, tstrerror(code));
808 809 810
  }
}

811
int32_t tmqHandleAllDelayedTask(tmq_t* pTmq) {
L
Liu Jicong 已提交
812
  STaosQall* qall = taosAllocateQall();
813
  taosReadAllQitems(pTmq->delayedTask, qall);
L
Liu Jicong 已提交
814

815 816 817 818
  if (qall->numOfItems == 0) {
    taosFreeQall(qall);
    return TSDB_CODE_SUCCESS;
  }
819

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

824
  while (pTaskType != NULL) {
825
    if (*pTaskType == TMQ_DELAYED_TASK__ASK_EP) {
826
      asyncAskEp(pTmq, addToQueueCallbackFn, NULL);
827 828

      int64_t* pRefId = taosMemoryMalloc(sizeof(int64_t));
829
      *pRefId = pTmq->refId;
830

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

      asyncCommitAllOffsets(pTmq, pCallbackFn, pTmq->commitCbUserParam);
837
      int64_t* pRefId = taosMemoryMalloc(sizeof(int64_t));
838
      *pRefId = pTmq->refId;
839

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

L
Liu Jicong 已提交
846
    taosFreeQitem(pTaskType);
847
    taosGetQitem(qall, (void**)&pTaskType);
L
Liu Jicong 已提交
848
  }
849

L
Liu Jicong 已提交
850 851 852 853
  taosFreeQall(qall);
  return 0;
}

854
static void* tmqFreeRspWrapper(SMqRspWrapper* rspWrapper) {
L
Liu Jicong 已提交
855 856 857 858 859 860 861
  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;
862 863
    taosMemoryFreeClear(pRsp->pEpset);

L
Liu Jicong 已提交
864 865 866 867 868 869
    taosArrayDestroyP(pRsp->dataRsp.blockData, taosMemoryFree);
    taosArrayDestroy(pRsp->dataRsp.blockDataLen);
    taosArrayDestroyP(pRsp->dataRsp.blockTbName, taosMemoryFree);
    taosArrayDestroyP(pRsp->dataRsp.blockSchema, (FDelete)tDeleteSSchemaWrapper);
  } else if (rspWrapper->tmqRspType == TMQ_MSG_TYPE__POLL_META_RSP) {
    SMqPollRspWrapper* pRsp = (SMqPollRspWrapper*)rspWrapper;
870 871
    taosMemoryFreeClear(pRsp->pEpset);

L
Liu Jicong 已提交
872 873 874
    taosMemoryFree(pRsp->metaRsp.metaRsp);
  } else if (rspWrapper->tmqRspType == TMQ_MSG_TYPE__TAOSX_RSP) {
    SMqPollRspWrapper* pRsp = (SMqPollRspWrapper*)rspWrapper;
875 876
    taosMemoryFreeClear(pRsp->pEpset);

L
Liu Jicong 已提交
877 878 879 880 881 882 883 884
    taosArrayDestroyP(pRsp->taosxRsp.blockData, taosMemoryFree);
    taosArrayDestroy(pRsp->taosxRsp.blockDataLen);
    taosArrayDestroyP(pRsp->taosxRsp.blockTbName, taosMemoryFree);
    taosArrayDestroyP(pRsp->taosxRsp.blockSchema, (FDelete)tDeleteSSchemaWrapper);
    // taosx
    taosArrayDestroy(pRsp->taosxRsp.createTableLen);
    taosArrayDestroyP(pRsp->taosxRsp.createTableReq, taosMemoryFree);
  }
885 886

  return NULL;
L
Liu Jicong 已提交
887 888
}

L
Liu Jicong 已提交
889
void tmqClearUnhandleMsg(tmq_t* tmq) {
L
Liu Jicong 已提交
890
  SMqRspWrapper* rspWrapper = NULL;
L
Liu Jicong 已提交
891
  while (1) {
L
Liu Jicong 已提交
892 893 894 895 896
    taosGetQitem(tmq->qall, (void**)&rspWrapper);
    if (rspWrapper) {
      tmqFreeRspWrapper(rspWrapper);
      taosFreeQitem(rspWrapper);
    } else {
L
Liu Jicong 已提交
897
      break;
L
Liu Jicong 已提交
898
    }
L
Liu Jicong 已提交
899 900
  }

L
Liu Jicong 已提交
901
  rspWrapper = NULL;
L
Liu Jicong 已提交
902 903
  taosReadAllQitems(tmq->mqueue, tmq->qall);
  while (1) {
L
Liu Jicong 已提交
904 905 906 907 908
    taosGetQitem(tmq->qall, (void**)&rspWrapper);
    if (rspWrapper) {
      tmqFreeRspWrapper(rspWrapper);
      taosFreeQitem(rspWrapper);
    } else {
L
Liu Jicong 已提交
909
      break;
L
Liu Jicong 已提交
910
    }
L
Liu Jicong 已提交
911 912 913
  }
}

D
dapan1121 已提交
914
int32_t tmqSubscribeCb(void* param, SDataBuf* pMsg, int32_t code) {
L
Liu Jicong 已提交
915 916
  SMqSubscribeCbParam* pParam = (SMqSubscribeCbParam*)param;
  pParam->rspErr = code;
dengyihao's avatar
dengyihao 已提交
917 918

  taosMemoryFree(pMsg->pEpSet);
L
Liu Jicong 已提交
919 920 921
  tsem_post(&pParam->rspSem);
  return 0;
}
922

L
Liu Jicong 已提交
923
int32_t tmq_subscription(tmq_t* tmq, tmq_list_t** topics) {
X
Xiaoyu Wang 已提交
924 925 926 927
  if (*topics == NULL) {
    *topics = tmq_list_new();
  }
  for (int i = 0; i < taosArrayGetSize(tmq->clientTopics); i++) {
L
Liu Jicong 已提交
928
    SMqClientTopic* topic = taosArrayGet(tmq->clientTopics, i);
L
Liu Jicong 已提交
929
    tmq_list_append(*topics, strchr(topic->topicName, '.') + 1);
X
Xiaoyu Wang 已提交
930
  }
L
Liu Jicong 已提交
931
  return 0;
X
Xiaoyu Wang 已提交
932 933
}

L
Liu Jicong 已提交
934
int32_t tmq_unsubscribe(tmq_t* tmq) {
L
Liu Jicong 已提交
935 936
  int32_t     rsp;
  int32_t     retryCnt = 0;
L
Liu Jicong 已提交
937
  tmq_list_t* lst = tmq_list_new();
L
Liu Jicong 已提交
938 939 940 941 942 943 944 945 946 947
  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 已提交
948 949
  tmq_list_destroy(lst);
  return rsp;
X
Xiaoyu Wang 已提交
950 951
}

952 953 954 955 956 957
static void freeClientVgImpl(void* param) {
  SMqClientTopic* pTopic = param;
  taosMemoryFreeClear(pTopic->schema.pSchema);
  taosArrayDestroy(pTopic->vgs);
}

958
void tmqFreeImpl(void* handle) {
959 960
  tmq_t*  tmq = (tmq_t*)handle;
  int64_t id = tmq->consumerId;
L
Liu Jicong 已提交
961

962
  // TODO stop timer
L
Liu Jicong 已提交
963 964 965 966
  if (tmq->mqueue) {
    tmqClearUnhandleMsg(tmq);
    taosCloseQueue(tmq->mqueue);
  }
L
Liu Jicong 已提交
967

H
Haojun Liao 已提交
968 969 970 971 972
  if (tmq->delayedTask) {
    taosCloseQueue(tmq->delayedTask);
  }

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

975
  taosArrayDestroyEx(tmq->clientTopics, freeClientVgImpl);
976 977
  taos_close_internal(tmq->pTscObj);
  taosMemoryFree(tmq);
978 979

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

982 983 984 985 986 987 988 989 990
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);
991
  if (tmqMgmt.rsetId < 0) {
992 993 994 995
    tmqInitRes = terrno;
  }
}

L
Liu Jicong 已提交
996
tmq_t* tmq_consumer_new(tmq_conf_t* conf, char* errstr, int32_t errstrLen) {
997 998 999 1000
  taosThreadOnce(&tmqInit, tmqMgmtInit);
  if (tmqInitRes != 0) {
    terrno = tmqInitRes;
    return NULL;
L
Liu Jicong 已提交
1001 1002
  }

L
Liu Jicong 已提交
1003 1004
  tmq_t* pTmq = taosMemoryCalloc(1, sizeof(tmq_t));
  if (pTmq == NULL) {
L
Liu Jicong 已提交
1005
    terrno = TSDB_CODE_OUT_OF_MEMORY;
1006
    tscError("failed to create consumer, groupId:%s, code:%s", conf->groupId, terrstr());
L
Liu Jicong 已提交
1007 1008
    return NULL;
  }
L
Liu Jicong 已提交
1009

L
Liu Jicong 已提交
1010 1011 1012
  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 已提交
1013 1014 1015
  pTmq->clientTopics = taosArrayInit(0, sizeof(SMqClientTopic));
  pTmq->mqueue = taosOpenQueue();
  pTmq->delayedTask = taosOpenQueue();
H
Haojun Liao 已提交
1016
  pTmq->qall = taosAllocateQall();
L
Liu Jicong 已提交
1017

X
Xiaoyu Wang 已提交
1018 1019
  if (pTmq->clientTopics == NULL || pTmq->mqueue == NULL || pTmq->qall == NULL || pTmq->delayedTask == NULL ||
      conf->groupId[0] == 0) {
L
Liu Jicong 已提交
1020
    terrno = TSDB_CODE_OUT_OF_MEMORY;
X
Xiaoyu Wang 已提交
1021
    tscError("consumer:0x%" PRIx64 " setup failed since %s, groupId:%s", pTmq->consumerId, terrstr(), pTmq->groupId);
1022
    goto _failed;
L
Liu Jicong 已提交
1023
  }
L
Liu Jicong 已提交
1024

L
Liu Jicong 已提交
1025 1026
  // init status
  pTmq->status = TMQ_CONSUMER_STATUS__INIT;
L
Liu Jicong 已提交
1027 1028
  pTmq->pollCnt = 0;
  pTmq->epoch = 0;
L
Liu Jicong 已提交
1029

L
Liu Jicong 已提交
1030 1031 1032
  // set conf
  strcpy(pTmq->clientId, conf->clientId);
  strcpy(pTmq->groupId, conf->groupId);
1033
  pTmq->withTbName = conf->withTbName;
L
Liu Jicong 已提交
1034
  pTmq->useSnapshot = conf->snapEnable;
L
Liu Jicong 已提交
1035
  pTmq->autoCommit = conf->autoCommit;
L
Liu Jicong 已提交
1036
  pTmq->autoCommitInterval = conf->autoCommitInterval;
L
Liu Jicong 已提交
1037 1038
  pTmq->commitCb = conf->commitCb;
  pTmq->commitCbUserParam = conf->commitCbUserParam;
L
Liu Jicong 已提交
1039 1040
  pTmq->resetOffsetCfg = conf->resetOffset;

1041 1042
  pTmq->hbBgEnable = conf->hbBgEnable;

L
Liu Jicong 已提交
1043
  // assign consumerId
L
Liu Jicong 已提交
1044
  pTmq->consumerId = tGenIdPI64();
X
Xiaoyu Wang 已提交
1045

L
Liu Jicong 已提交
1046 1047
  // init semaphore
  if (tsem_init(&pTmq->rspSem, 0, 0) != 0) {
1048
    tscError("consumer:0x %" PRIx64 " setup failed since %s, consumer group %s", pTmq->consumerId, terrstr(),
S
Shengliang Guan 已提交
1049
             pTmq->groupId);
1050
    goto _failed;
L
Liu Jicong 已提交
1051
  }
L
Liu Jicong 已提交
1052

L
Liu Jicong 已提交
1053 1054 1055
  // init connection
  pTmq->pTscObj = taos_connect_internal(conf->ip, user, pass, NULL, NULL, conf->port, CONN_TYPE__TMQ);
  if (pTmq->pTscObj == NULL) {
1056
    tscError("consumer:0x%" PRIx64 " setup failed since %s, groupId:%s", pTmq->consumerId, terrstr(), pTmq->groupId);
L
Liu Jicong 已提交
1057
    tsem_destroy(&pTmq->rspSem);
1058
    goto _failed;
L
Liu Jicong 已提交
1059
  }
L
Liu Jicong 已提交
1060

1061 1062
  pTmq->refId = taosAddRef(tmqMgmt.rsetId, pTmq);
  if (pTmq->refId < 0) {
1063
    goto _failed;
1064 1065
  }

1066
  if (pTmq->hbBgEnable) {
L
Liu Jicong 已提交
1067 1068
    int64_t* pRefId = taosMemoryMalloc(sizeof(int64_t));
    *pRefId = pTmq->refId;
1069
    pTmq->hbLiveTimer = taosTmrStart(tmqSendHbReq, 1000, pRefId, tmqMgmt.timer);
1070 1071
  }

1072 1073 1074
  char         buf[80] = {0};
  STqOffsetVal offset = {.type = pTmq->resetOffsetCfg};
  tFormatOffset(buf, tListLen(buf), &offset);
X
Xiaoyu Wang 已提交
1075 1076 1077 1078
  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 已提交
1079

1080
  return pTmq;
1081

1082 1083
_failed:
  tmqFreeImpl(pTmq);
L
Liu Jicong 已提交
1084
  return NULL;
1085 1086
}

L
Liu Jicong 已提交
1087
int32_t tmq_subscribe(tmq_t* tmq, const tmq_list_t* topic_list) {
1088
  const int32_t   MAX_RETRY_COUNT = 120 * 2;  // let's wait for 2 mins at most
L
Liu Jicong 已提交
1089 1090 1091
  const SArray*   container = &topic_list->container;
  int32_t         sz = taosArrayGetSize(container);
  void*           buf = NULL;
L
Liu Jicong 已提交
1092
  SMsgSendInfo*   sendInfo = NULL;
L
Liu Jicong 已提交
1093
  SCMSubscribeReq req = {0};
1094
  int32_t         code = 0;
1095

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

1098
  req.consumerId = tmq->consumerId;
L
Liu Jicong 已提交
1099
  tstrncpy(req.clientId, tmq->clientId, 256);
L
Liu Jicong 已提交
1100
  tstrncpy(req.cgroup, tmq->groupId, TSDB_CGROUP_LEN);
1101 1102
  req.topicNames = taosArrayInit(sz, sizeof(void*));

1103 1104 1105 1106
  if (req.topicNames == NULL) {
    code = TSDB_CODE_OUT_OF_MEMORY;
    goto FAIL;
  }
1107

L
Liu Jicong 已提交
1108 1109
  for (int32_t i = 0; i < sz; i++) {
    char* topic = taosArrayGetP(container, i);
1110 1111

    SName name = {0};
L
Liu Jicong 已提交
1112 1113 1114 1115
    tNameSetDbName(&name, tmq->pTscObj->acctId, topic, strlen(topic));
    char* topicFName = taosMemoryCalloc(1, TSDB_TOPIC_FNAME_LEN);
    if (topicFName == NULL) {
      goto FAIL;
1116 1117
    }

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

    taosArrayPush(req.topicNames, &topicFName);
1122 1123
  }

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

L
Liu Jicong 已提交
1126
  buf = taosMemoryMalloc(tlen);
1127 1128 1129 1130
  if (buf == NULL) {
    code = TSDB_CODE_OUT_OF_MEMORY;
    goto FAIL;
  }
L
Liu Jicong 已提交
1131

1132 1133 1134
  void* abuf = buf;
  tSerializeSCMSubscribeReq(&abuf, &req);

L
Liu Jicong 已提交
1135
  sendInfo = taosMemoryCalloc(1, sizeof(SMsgSendInfo));
1136 1137 1138 1139
  if (sendInfo == NULL) {
    code = TSDB_CODE_OUT_OF_MEMORY;
    goto FAIL;
  }
1140

X
Xiaoyu Wang 已提交
1141
  SMqSubscribeCbParam param = {
L
Liu Jicong 已提交
1142
      .rspErr = 0,
1143 1144
      .refId = tmq->refId,
      .epoch = tmq->epoch,
X
Xiaoyu Wang 已提交
1145
  };
L
Liu Jicong 已提交
1146

1147
  if (tsem_init(&param.rspSem, 0, 0) != 0) {
wmmhello's avatar
wmmhello 已提交
1148
    code = TSDB_CODE_TSC_INTERNAL_ERROR;
1149 1150
    goto FAIL;
  }
L
Liu Jicong 已提交
1151 1152

  sendInfo->msgInfo = (SDataBuf){
X
Xiaoyu Wang 已提交
1153 1154 1155 1156
      .pData = buf,
      .len = tlen,
      .handle = NULL,
  };
1157

L
Liu Jicong 已提交
1158 1159
  sendInfo->requestId = generateRequestId();
  sendInfo->requestObjRefId = 0;
L
Liu Jicong 已提交
1160 1161
  sendInfo->param = &param;
  sendInfo->fp = tmqSubscribeCb;
L
Liu Jicong 已提交
1162
  sendInfo->msgType = TDMT_MND_TMQ_SUBSCRIBE;
L
Liu Jicong 已提交
1163

1164 1165 1166 1167 1168
  SEpSet epSet = getEpSet_s(&tmq->pTscObj->pAppInfo->mgmtEp);

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

L
Liu Jicong 已提交
1169 1170
  // avoid double free if msg is sent
  buf = NULL;
L
Liu Jicong 已提交
1171
  sendInfo = NULL;
L
Liu Jicong 已提交
1172

L
Liu Jicong 已提交
1173 1174
  tsem_wait(&param.rspSem);
  tsem_destroy(&param.rspSem);
1175

1176 1177 1178 1179
  if (param.rspErr != 0) {
    code = param.rspErr;
    goto FAIL;
  }
L
Liu Jicong 已提交
1180

L
Liu Jicong 已提交
1181
  int32_t retryCnt = 0;
1182
  while (TSDB_CODE_MND_CONSUMER_NOT_READY == doAskEp(tmq)) {
1183
    if (retryCnt++ > MAX_RETRY_COUNT) {
wmmhello's avatar
wmmhello 已提交
1184 1185
      tscError("consumer:0x%" PRIx64 ", mnd not ready for subscribe, retry:%d in 500ms", tmq->consumerId, retryCnt);
      code = TSDB_CODE_TSC_INTERNAL_ERROR;
L
Liu Jicong 已提交
1186 1187
      goto FAIL;
    }
1188

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

1193 1194
  // init ep timer
  if (tmq->epTimer == NULL) {
1195 1196 1197
    int64_t* pRefId1 = taosMemoryMalloc(sizeof(int64_t));
    *pRefId1 = tmq->refId;
    tmq->epTimer = taosTmrStart(tmqAssignAskEpTask, 1000, pRefId1, tmqMgmt.timer);
1198
  }
L
Liu Jicong 已提交
1199 1200

  // init auto commit timer
1201
  if (tmq->autoCommit && tmq->commitTimer == NULL) {
1202 1203 1204
    int64_t* pRefId2 = taosMemoryMalloc(sizeof(int64_t));
    *pRefId2 = tmq->refId;
    tmq->commitTimer = taosTmrStart(tmqAssignDelayedCommitTask, tmq->autoCommitInterval, pRefId2, tmqMgmt.timer);
L
Liu Jicong 已提交
1205 1206
  }

L
Liu Jicong 已提交
1207
FAIL:
L
Liu Jicong 已提交
1208
  taosArrayDestroyP(req.topicNames, taosMemoryFree);
L
Liu Jicong 已提交
1209
  taosMemoryFree(buf);
L
Liu Jicong 已提交
1210
  taosMemoryFree(sendInfo);
L
Liu Jicong 已提交
1211

L
Liu Jicong 已提交
1212
  return code;
1213 1214
}

L
Liu Jicong 已提交
1215
void tmq_conf_set_auto_commit_cb(tmq_conf_t* conf, tmq_commit_cb* cb, void* param) {
1216
  conf->commitCb = cb;
L
Liu Jicong 已提交
1217
  conf->commitCbUserParam = param;
L
Liu Jicong 已提交
1218
}
1219

D
dapan1121 已提交
1220
int32_t tmqPollCb(void* param, SDataBuf* pMsg, int32_t code) {
X
Xiaoyu Wang 已提交
1221
  SMqPollCbParam* pParam = (SMqPollCbParam*)param;
1222 1223

  int64_t         refId = pParam->refId;
X
Xiaoyu Wang 已提交
1224
  SMqClientVg*    pVg = pParam->pVg;
L
Liu Jicong 已提交
1225
  SMqClientTopic* pTopic = pParam->pTopic;
1226

1227
  tmq_t* tmq = taosAcquireRef(tmqMgmt.rsetId, refId);
1228 1229 1230
  if (tmq == NULL) {
    taosMemoryFree(pParam);
    taosMemoryFree(pMsg->pData);
dengyihao's avatar
dengyihao 已提交
1231
    taosMemoryFree(pMsg->pEpSet);
1232 1233 1234 1235
    terrno = TSDB_CODE_TMQ_CONSUMER_CLOSED;
    return -1;
  }

H
Haojun Liao 已提交
1236 1237 1238 1239
  int32_t  epoch = pParam->epoch;
  int32_t  vgId = pParam->vgId;
  uint64_t requestId = pParam->requestId;

L
Liu Jicong 已提交
1240
  taosMemoryFree(pParam);
H
Haojun Liao 已提交
1241

L
Liu Jicong 已提交
1242
  if (code != 0) {
L
Liu Jicong 已提交
1243
    if (pMsg->pData) taosMemoryFree(pMsg->pData);
dengyihao's avatar
dengyihao 已提交
1244 1245
    if (pMsg->pEpSet) taosMemoryFree(pMsg->pEpSet);

H
Haojun Liao 已提交
1246
    // in case of consumer mismatch, wait for 500ms and retry
L
Liu Jicong 已提交
1247
    if (code == TSDB_CODE_TMQ_CONSUMER_MISMATCH) {
1248
      taosMsleep(500);
L
Liu Jicong 已提交
1249
      atomic_store_8(&tmq->status, TMQ_CONSUMER_STATUS__RECOVER);
X
Xiaoyu Wang 已提交
1250 1251
      tscDebug("consumer:0x%" PRIx64 " wait for the re-balance, wait for 500ms and set status to be RECOVER",
               tmq->consumerId);
H
Haojun Liao 已提交
1252
    } else if (code == TSDB_CODE_TQ_NO_COMMITTED_OFFSET) {
S
Shengliang Guan 已提交
1253
      SMqPollRspWrapper* pRspWrapper = taosAllocateQitem(sizeof(SMqPollRspWrapper), DEF_QITEM, 0);
L
Liu Jicong 已提交
1254
      if (pRspWrapper == NULL) {
H
Haojun Liao 已提交
1255 1256
        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 已提交
1257 1258
        goto CREATE_MSG_FAIL;
      }
H
Haojun Liao 已提交
1259

L
Liu Jicong 已提交
1260 1261
      pRspWrapper->tmqRspType = TMQ_MSG_TYPE__END_RSP;
      taosWriteQitem(tmq->mqueue, pRspWrapper);
X
Xiaoyu Wang 已提交
1262
    } else if (code == TSDB_CODE_WAL_LOG_NOT_EXIST) {  // poll data while insert
1263
      taosMsleep(500);
wmmhello's avatar
wmmhello 已提交
1264 1265 1266
    } 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 已提交
1267
    }
H
Haojun Liao 已提交
1268

L
fix txn  
Liu Jicong 已提交
1269
    goto CREATE_MSG_FAIL;
L
Liu Jicong 已提交
1270 1271
  }

X
Xiaoyu Wang 已提交
1272 1273 1274
  int32_t msgEpoch = ((SMqRspHead*)pMsg->pData)->epoch;
  int32_t tmqEpoch = atomic_load_32(&tmq->epoch);
  if (msgEpoch < tmqEpoch) {
L
Liu Jicong 已提交
1275
    // do not write into queue since updating epoch reset
X
Xiaoyu Wang 已提交
1276 1277
    tscWarn("consumer:0x%" PRIx64
            " msg discard from vgId:%d since from earlier epoch, rsp epoch %d, current epoch %d, reqId:0x%" PRIx64,
H
Haojun Liao 已提交
1278 1279
            tmq->consumerId, vgId, msgEpoch, tmqEpoch, requestId);

1280
    tsem_post(&tmq->rspSem);
1281 1282
    taosReleaseRef(tmqMgmt.rsetId, refId);

L
Liu Jicong 已提交
1283
    taosMemoryFree(pMsg->pData);
dengyihao's avatar
dengyihao 已提交
1284
    taosMemoryFree(pMsg->pEpSet);
X
Xiaoyu Wang 已提交
1285 1286 1287 1288
    return 0;
  }

  if (msgEpoch != tmqEpoch) {
H
Haojun Liao 已提交
1289 1290
    tscWarn("consumer:0x%" PRIx64 " mismatch rsp from vgId:%d, epoch %d, current epoch %d, reqId:0x%" PRIx64,
            tmq->consumerId, vgId, msgEpoch, tmqEpoch, requestId);
X
Xiaoyu Wang 已提交
1291 1292
  }

L
Liu Jicong 已提交
1293 1294 1295
  // handle meta rsp
  int8_t rspType = ((SMqRspHead*)pMsg->pData)->mqMsgType;

S
Shengliang Guan 已提交
1296
  SMqPollRspWrapper* pRspWrapper = taosAllocateQitem(sizeof(SMqPollRspWrapper), DEF_QITEM, 0);
L
Liu Jicong 已提交
1297
  if (pRspWrapper == NULL) {
L
Liu Jicong 已提交
1298
    taosMemoryFree(pMsg->pData);
dengyihao's avatar
dengyihao 已提交
1299
    taosMemoryFree(pMsg->pEpSet);
X
Xiaoyu Wang 已提交
1300 1301
    tscWarn("consumer:0x%" PRIx64 " msg discard from vgId:%d, epoch %d since out of memory", tmq->consumerId, vgId,
            epoch);
L
fix txn  
Liu Jicong 已提交
1302
    goto CREATE_MSG_FAIL;
L
Liu Jicong 已提交
1303
  }
L
Liu Jicong 已提交
1304

L
Liu Jicong 已提交
1305
  pRspWrapper->tmqRspType = rspType;
L
Liu Jicong 已提交
1306 1307
  pRspWrapper->vgHandle = pVg;
  pRspWrapper->topicHandle = pTopic;
H
Haojun Liao 已提交
1308
  pRspWrapper->reqId = requestId;
1309
  pRspWrapper->pEpset = pMsg->pEpSet;
1310
  pRspWrapper->vgId = pVg->vgId;
L
Liu Jicong 已提交
1311

1312
  pMsg->pEpSet = NULL;
L
Liu Jicong 已提交
1313
  if (rspType == TMQ_MSG_TYPE__POLL_RSP) {
L
Liu Jicong 已提交
1314 1315 1316
    SDecoder decoder;
    tDecoderInit(&decoder, POINTER_SHIFT(pMsg->pData, sizeof(SMqRspHead)), pMsg->len - sizeof(SMqRspHead));
    tDecodeSMqDataRsp(&decoder, &pRspWrapper->dataRsp);
wmmhello's avatar
wmmhello 已提交
1317
    tDecoderClear(&decoder);
L
Liu Jicong 已提交
1318
    memcpy(&pRspWrapper->dataRsp, pMsg->pData, sizeof(SMqRspHead));
L
Liu Jicong 已提交
1319

H
Haojun Liao 已提交
1320 1321
    char buf[80];
    tFormatOffset(buf, 80, &pRspWrapper->dataRsp.rspOffset);
H
Haojun Liao 已提交
1322
    tscDebug("consumer:0x%" PRIx64 " recv poll rsp, vgId:%d, req ver:%" PRId64 ", rsp:%s type %d, reqId:0x%" PRIx64,
H
Haojun Liao 已提交
1323
             tmq->consumerId, vgId, pRspWrapper->dataRsp.reqOffset.version, buf, rspType, requestId);
L
Liu Jicong 已提交
1324
  } else if (rspType == TMQ_MSG_TYPE__POLL_META_RSP) {
1325 1326 1327 1328
    SDecoder decoder;
    tDecoderInit(&decoder, POINTER_SHIFT(pMsg->pData, sizeof(SMqRspHead)), pMsg->len - sizeof(SMqRspHead));
    tDecodeSMqMetaRsp(&decoder, &pRspWrapper->metaRsp);
    tDecoderClear(&decoder);
L
Liu Jicong 已提交
1329
    memcpy(&pRspWrapper->metaRsp, pMsg->pData, sizeof(SMqRspHead));
L
Liu Jicong 已提交
1330 1331 1332 1333 1334 1335
  } 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 已提交
1336 1337
  } else {  // invalid rspType
    tscError("consumer:0x%" PRIx64 " invalid rsp msg received, type:%d ignored", tmq->consumerId, rspType);
L
Liu Jicong 已提交
1338
  }
L
Liu Jicong 已提交
1339

L
Liu Jicong 已提交
1340
  taosMemoryFree(pMsg->pData);
H
Haojun Liao 已提交
1341
  taosWriteQitem(tmq->mqueue, pRspWrapper);
L
Liu Jicong 已提交
1342

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

1347
  tsem_post(&tmq->rspSem);
1348 1349
  taosReleaseRef(tmqMgmt.rsetId, refId);

L
Liu Jicong 已提交
1350
  return 0;
H
Haojun Liao 已提交
1351

L
fix txn  
Liu Jicong 已提交
1352
CREATE_MSG_FAIL:
L
Liu Jicong 已提交
1353
  if (epoch == tmq->epoch) {
L
Liu Jicong 已提交
1354 1355
    atomic_store_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE);
  }
H
Haojun Liao 已提交
1356

1357
  tsem_post(&tmq->rspSem);
1358 1359
  taosReleaseRef(tmqMgmt.rsetId, refId);

L
Liu Jicong 已提交
1360
  return -1;
1361 1362
}

H
Haojun Liao 已提交
1363 1364 1365 1366 1367
typedef struct SVgroupSaveInfo {
  STqOffsetVal offset;
  int64_t      numOfRows;
} SVgroupSaveInfo;

H
Haojun Liao 已提交
1368 1369 1370 1371 1372 1373
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 已提交
1374
  char    vgKey[TSDB_TOPIC_FNAME_LEN + 22];
H
Haojun Liao 已提交
1375 1376 1377 1378 1379
  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 已提交
1380
  tscDebug("consumer:0x%" PRIx64 ", update topic:%s, new numOfVgs:%d", tmq->consumerId, pTopic->topicName, vgNumGet);
H
Haojun Liao 已提交
1381 1382 1383 1384
  pTopic->vgs = taosArrayInit(vgNumGet, sizeof(SMqClientVg));

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

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

X
Xiaoyu Wang 已提交
1389 1390
    int64_t      numOfRows = 0;
    STqOffsetVal offsetNew = {.type = tmq->resetOffsetCfg};
H
Haojun Liao 已提交
1391 1392 1393
    if (pInfo != NULL) {
      offsetNew = pInfo->offset;
      numOfRows = pInfo->numOfRows;
H
Haojun Liao 已提交
1394 1395 1396 1397 1398 1399 1400
    }

    SMqClientVg clientVg = {
        .pollCnt = 0,
        .currentOffset = offsetNew,
        .vgId = pVgEp->vgId,
        .epSet = pVgEp->epSet,
wmmhello's avatar
wmmhello 已提交
1401
        .vgStatus = TMQ_VG_STATUS__IDLE,
H
Haojun Liao 已提交
1402
        .vgSkipCnt = 0,
H
Haojun Liao 已提交
1403
        .emptyBlockReceiveTs = 0,
H
Haojun Liao 已提交
1404
        .numOfRows = numOfRows,
H
Haojun Liao 已提交
1405 1406 1407 1408 1409 1410 1411 1412 1413 1414 1415 1416 1417 1418 1419
    };

    taosArrayPush(pTopic->vgs, &clientVg);
  }
}

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

  taosArrayDestroy(pTopic->vgs);
}

1420
static bool doUpdateLocalEp(tmq_t* tmq, int32_t epoch, const SMqAskEpRsp* pRsp) {
1421 1422
  bool set = false;

1423
  int32_t topicNumCur = taosArrayGetSize(tmq->clientTopics);
1424
  int32_t topicNumGet = taosArrayGetSize(pRsp->topics);
1425

X
Xiaoyu Wang 已提交
1426 1427
  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",
1428
           tmq->consumerId, tmq->epoch, epoch, topicNumGet, topicNumCur);
wmmhello's avatar
wmmhello 已提交
1429 1430 1431
  if (epoch <= tmq->epoch) {
    return false;
  }
1432 1433 1434 1435 1436 1437

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

H
Haojun Liao 已提交
1438 1439
  SHashObj* pVgOffsetHashMap = taosHashInit(64, MurmurHash3_32, false, HASH_NO_LOCK);
  if (pVgOffsetHashMap == NULL) {
1440 1441 1442
    taosArrayDestroy(newTopics);
    return false;
  }
1443

H
Haojun Liao 已提交
1444
  // todo extract method
1445 1446 1447 1448 1449
  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 已提交
1450
      tscDebug("consumer:0x%" PRIx64 ", current vg num: %d", tmq->consumerId, vgNumCur);
1451 1452
      for (int32_t j = 0; j < vgNumCur; j++) {
        SMqClientVg* pVgCur = taosArrayGet(pTopicCur->vgs, j);
H
Haojun Liao 已提交
1453 1454
        makeTopicVgroupKey(vgKey, pTopicCur->topicName, pVgCur->vgId);

L
Liu Jicong 已提交
1455
        char buf[80];
L
Liu Jicong 已提交
1456
        tFormatOffset(buf, 80, &pVgCur->currentOffset);
wmmhello's avatar
wmmhello 已提交
1457
        tscDebug("consumer:0x%" PRIx64 ", doUpdateLocalEp current vg, epoch:%d vgId:%d vgKey:%s, offset:%s", tmq->consumerId, tmq->epoch, pVgCur->vgId,
X
Xiaoyu Wang 已提交
1458
                 vgKey, buf);
H
Haojun Liao 已提交
1459

wmmhello's avatar
wmmhello 已提交
1460
        SVgroupSaveInfo info = {.offset = pVgCur->currentOffset, .numOfRows = pVgCur->numOfRows};
H
Haojun Liao 已提交
1461
        taosHashPut(pVgOffsetHashMap, vgKey, strlen(vgKey), &info, sizeof(SVgroupSaveInfo));
1462 1463 1464 1465 1466 1467 1468
      }
    }
  }

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

H
Haojun Liao 已提交
1473 1474
  taosHashCleanup(pVgOffsetHashMap);

1475
  // destroy current buffered existed topics info
1476
  if (tmq->clientTopics) {
H
Haojun Liao 已提交
1477
    taosArrayDestroyEx(tmq->clientTopics, freeClientVgInfo);
X
Xiaoyu Wang 已提交
1478
  }
H
Haojun Liao 已提交
1479
  tmq->clientTopics = newTopics;
1480

X
Xiaoyu Wang 已提交
1481
  int8_t flag = (topicNumGet == 0) ? TMQ_CONSUMER_STATUS__NO_TOPIC : TMQ_CONSUMER_STATUS__READY;
H
Haojun Liao 已提交
1482
  atomic_store_8(&tmq->status, flag);
X
Xiaoyu Wang 已提交
1483
  atomic_store_32(&tmq->epoch, epoch);
H
Haojun Liao 已提交
1484

1485
  tscDebug("consumer:0x%" PRIx64 " update topic info completed", tmq->consumerId);
X
Xiaoyu Wang 已提交
1486 1487 1488
  return set;
}

1489
int32_t askEpCallbackFn(void* param, SDataBuf* pMsg, int32_t code) {
1490
  SMqAskEpCbParam* pParam = (SMqAskEpCbParam*)param;
1491 1492 1493
  tmq_t*           tmq = taosAcquireRef(tmqMgmt.rsetId, pParam->refId);

  if (tmq == NULL) {
1494 1495 1496
    terrno = TSDB_CODE_TMQ_CONSUMER_CLOSED;
    pParam->pUserFn(tmq, terrno, NULL, pParam->pParam);

1497
    taosMemoryFree(pMsg->pData);
dengyihao's avatar
dengyihao 已提交
1498
    taosMemoryFree(pMsg->pEpSet);
1499 1500
    taosMemoryFree(pParam);
    return terrno;
1501 1502
  }

H
Haojun Liao 已提交
1503
  if (code != TSDB_CODE_SUCCESS) {
1504 1505 1506 1507 1508 1509 1510 1511 1512
    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;
1513
  }
L
Liu Jicong 已提交
1514

L
Liu Jicong 已提交
1515
  // tmq's epoch is monotonically increase,
L
Liu Jicong 已提交
1516
  // so it's safe to discard any old epoch msg.
L
Liu Jicong 已提交
1517
  // Epoch will only increase when received newer epoch ep msg
L
Liu Jicong 已提交
1518 1519 1520
  SMqRspHead* head = pMsg->pData;
  int32_t     epoch = atomic_load_32(&tmq->epoch);
  if (head->epoch <= epoch) {
1521 1522
    tscDebug("consumer:0x%" PRIx64 ", recv ep, msg epoch %d, current epoch %d, no need to update local ep",
             tmq->consumerId, head->epoch, epoch);
1523

1524 1525 1526 1527 1528 1529 1530 1531
    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 已提交
1532
  } else {
1533 1534
    tscDebug("consumer:0x%" PRIx64 ", recv ep, msg epoch %d, current epoch %d, update local ep", tmq->consumerId,
             head->epoch, epoch);
1535
  }
wmmhello's avatar
wmmhello 已提交
1536
  pParam->pUserFn(tmq, code, pMsg, pParam->pParam);
L
Liu Jicong 已提交
1537

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

dengyihao's avatar
dengyihao 已提交
1540
  taosMemoryFree(pMsg->pEpSet);
L
Liu Jicong 已提交
1541
  taosMemoryFree(pMsg->pData);
1542
  taosMemoryFree(pParam);
L
Liu Jicong 已提交
1543
  return code;
1544 1545
}

L
Liu Jicong 已提交
1546
void tmqBuildConsumeReqImpl(SMqPollReq* pReq, tmq_t* tmq, int64_t timeout, SMqClientTopic* pTopic, SMqClientVg* pVg) {
L
Liu Jicong 已提交
1547 1548 1549 1550
  int32_t groupLen = strlen(tmq->groupId);
  memcpy(pReq->subKey, tmq->groupId, groupLen);
  pReq->subKey[groupLen] = TMQ_SEPARATOR;
  strcpy(pReq->subKey + groupLen + 1, pTopic->topicName);
1551

1552
  pReq->withTbName = tmq->withTbName;
L
Liu Jicong 已提交
1553
  pReq->consumerId = tmq->consumerId;
1554
  pReq->timeout = timeout;
X
Xiaoyu Wang 已提交
1555
  pReq->epoch = tmq->epoch;
L
Liu Jicong 已提交
1556
  /*pReq->currentOffset = reqOffset;*/
L
Liu Jicong 已提交
1557
  pReq->reqOffset = pVg->currentOffset;
D
dapan1121 已提交
1558
  pReq->head.vgId = pVg->vgId;
1559 1560
  pReq->useSnapshot = tmq->useSnapshot;
  pReq->reqId = generateRequestId();
1561 1562
}

L
Liu Jicong 已提交
1563 1564
SMqMetaRspObj* tmqBuildMetaRspFromWrapper(SMqPollRspWrapper* pWrapper) {
  SMqMetaRspObj* pRspObj = taosMemoryCalloc(1, sizeof(SMqMetaRspObj));
L
Liu Jicong 已提交
1565
  pRspObj->resType = RES_TYPE__TMQ_META;
L
Liu Jicong 已提交
1566 1567 1568 1569 1570 1571 1572 1573
  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;
}

1574
SMqRspObj* tmqBuildRspFromWrapper(SMqPollRspWrapper* pWrapper, SMqClientVg* pVg, int64_t* numOfRows) {
L
Liu Jicong 已提交
1575 1576
  SMqRspObj* pRspObj = taosMemoryCalloc(1, sizeof(SMqRspObj));
  pRspObj->resType = RES_TYPE__TMQ;
1577

1578
  (*numOfRows) = 0;
L
Liu Jicong 已提交
1579 1580
  tstrncpy(pRspObj->topic, pWrapper->topicHandle->topicName, TSDB_TOPIC_FNAME_LEN);
  tstrncpy(pRspObj->db, pWrapper->topicHandle->db, TSDB_DB_FNAME_LEN);
1581

L
Liu Jicong 已提交
1582
  pRspObj->vgId = pWrapper->vgHandle->vgId;
L
Liu Jicong 已提交
1583
  pRspObj->resIter = -1;
L
Liu Jicong 已提交
1584
  memcpy(&pRspObj->rsp, &pWrapper->dataRsp, sizeof(SMqDataRsp));
L
Liu Jicong 已提交
1585

L
Liu Jicong 已提交
1586 1587
  pRspObj->resInfo.totalRows = 0;
  pRspObj->resInfo.precision = TSDB_TIME_PRECISION_MILLI;
1588

L
Liu Jicong 已提交
1589
  if (!pWrapper->dataRsp.withSchema) {
L
Liu Jicong 已提交
1590 1591
    setResSchemaInfo(&pRspObj->resInfo, pWrapper->topicHandle->schema.pSchema, pWrapper->topicHandle->schema.nCols);
  }
L
Liu Jicong 已提交
1592

1593
  // extract the rows in this data packet
X
Xiaoyu Wang 已提交
1594
  for (int32_t i = 0; i < pRspObj->rsp.blockNum; ++i) {
1595
    SRetrieveTableRsp* pRetrieve = (SRetrieveTableRsp*)taosArrayGetP(pRspObj->rsp.blockData, i);
X
Xiaoyu Wang 已提交
1596
    int64_t            rows = htobe64(pRetrieve->numOfRows);
1597
    pVg->numOfRows += rows;
1598
    (*numOfRows) += rows;
1599 1600
  }

L
Liu Jicong 已提交
1601
  return pRspObj;
X
Xiaoyu Wang 已提交
1602 1603
}

L
Liu Jicong 已提交
1604 1605
SMqTaosxRspObj* tmqBuildTaosxRspFromWrapper(SMqPollRspWrapper* pWrapper) {
  SMqTaosxRspObj* pRspObj = taosMemoryCalloc(1, sizeof(SMqTaosxRspObj));
1606
  pRspObj->resType = RES_TYPE__TMQ_METADATA;
L
Liu Jicong 已提交
1607 1608 1609 1610
  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;
1611
  memcpy(&pRspObj->rsp, &pWrapper->taosxRsp, sizeof(STaosxRsp));
L
Liu Jicong 已提交
1612 1613 1614

  pRspObj->resInfo.totalRows = 0;
  pRspObj->resInfo.precision = TSDB_TIME_PRECISION_MILLI;
1615
  if (!pWrapper->taosxRsp.withSchema) {
L
Liu Jicong 已提交
1616 1617 1618 1619 1620 1621
    setResSchemaInfo(&pRspObj->resInfo, pWrapper->topicHandle->schema.pSchema, pWrapper->topicHandle->schema.nCols);
  }

  return pRspObj;
}

1622 1623 1624 1625 1626 1627 1628 1629 1630 1631 1632 1633 1634 1635 1636 1637 1638 1639 1640 1641 1642 1643 1644 1645 1646 1647 1648 1649 1650 1651 1652 1653 1654
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;
X
Xiaoyu Wang 已提交
1655
  pParam->pVg = pVg;  // pVg may be released,fix it
1656 1657
  pParam->pTopic = pTopic;
  pParam->vgId = pVg->vgId;
H
Haojun Liao 已提交
1658
  pParam->requestId = req.reqId;
1659 1660 1661 1662 1663 1664 1665 1666

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

H
Haojun Liao 已提交
1667
  sendInfo->msgInfo = (SDataBuf){ .pData = msg, .len = msgSize, .handle = NULL };
1668 1669 1670 1671 1672 1673 1674 1675 1676 1677 1678

  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];
  tFormatOffset(offsetFormatBuf, tListLen(offsetFormatBuf), &pVg->currentOffset);

X
Xiaoyu Wang 已提交
1679 1680
  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);
1681 1682 1683 1684 1685 1686 1687 1688
  asyncSendMsgToServer(pTmq->pTscObj->pAppInfo->pTransporter, &pVg->epSet, &transporterId, sendInfo);

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

  return TSDB_CODE_SUCCESS;
}

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

  for (int i = 0; i < numOfTopics; i++) {
X
Xiaoyu Wang 已提交
1695
    SMqClientTopic* pTopic = taosArrayGet(tmq->clientTopics, i);
X
Xiaoyu Wang 已提交
1696
    int32_t         numOfVg = taosArrayGetSize(pTopic->vgs);
1697 1698

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

1706
      int32_t vgStatus = atomic_val_compare_exchange_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE, TMQ_VG_STATUS__WAIT);
1707
      if (vgStatus == TMQ_VG_STATUS__WAIT) {
L
Liu Jicong 已提交
1708
        int32_t vgSkipCnt = atomic_add_fetch_32(&pVg->vgSkipCnt, 1);
1709
        tscTrace("consumer:0x%" PRIx64 " epoch %d wait poll-rsp, skip vgId:%d skip cnt %d", tmq->consumerId, tmq->epoch,
X
Xiaoyu Wang 已提交
1710
                 pVg->vgId, vgSkipCnt);
X
Xiaoyu Wang 已提交
1711
        continue;
L
temp  
Liu Jicong 已提交
1712 1713 1714 1715
#if 0
        if (skipCnt < 30000) {
          continue;
        } else {
1716
        tscDebug("consumer:0x%" PRIx64 ",skip vgId:%d skip too much reset", tmq->consumerId, pVg->vgId);
L
temp  
Liu Jicong 已提交
1717 1718
        }
#endif
X
Xiaoyu Wang 已提交
1719
      }
1720

L
Liu Jicong 已提交
1721
      atomic_store_32(&pVg->vgSkipCnt, 0);
1722 1723 1724
      int32_t code = doTmqPollImpl(tmq, pTopic, pVg, timeout);
      if (code != TSDB_CODE_SUCCESS) {
        return code;
D
dapan1121 已提交
1725
      }
X
Xiaoyu Wang 已提交
1726 1727
    }
  }
1728

1729
  tscDebug("consumer:0x%" PRIx64 " end to poll data", tmq->consumerId);
X
Xiaoyu Wang 已提交
1730 1731 1732
  return 0;
}

H
Haojun Liao 已提交
1733
static int32_t tmqHandleNoPollRsp(tmq_t* tmq, SMqRspWrapper* rspWrapper, bool* pReset) {
L
Liu Jicong 已提交
1734
  if (rspWrapper->tmqRspType == TMQ_MSG_TYPE__EP_RSP) {
L
fix  
Liu Jicong 已提交
1735
    /*printf("ep %d %d\n", rspMsg->head.epoch, tmq->epoch);*/
L
Liu Jicong 已提交
1736 1737
    if (rspWrapper->epoch > atomic_load_32(&tmq->epoch)) {
      SMqAskEpRspWrapper* pEpRspWrapper = (SMqAskEpRspWrapper*)rspWrapper;
L
Liu Jicong 已提交
1738
      SMqAskEpRsp*        rspMsg = &pEpRspWrapper->msg;
1739
      doUpdateLocalEp(tmq, rspWrapper->epoch, rspMsg);
L
temp  
Liu Jicong 已提交
1740
      /*tmqClearUnhandleMsg(tmq);*/
L
Liu Jicong 已提交
1741
      tDeleteSMqAskEpRsp(rspMsg);
X
Xiaoyu Wang 已提交
1742 1743
      *pReset = true;
    } else {
L
Liu Jicong 已提交
1744
      tmqFreeRspWrapper(rspWrapper);
X
Xiaoyu Wang 已提交
1745 1746 1747 1748 1749 1750 1751 1752
      *pReset = false;
    }
  } else {
    return -1;
  }
  return 0;
}

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

X
Xiaoyu Wang 已提交
1756
  while (1) {
1757 1758
    SMqRspWrapper* pRspWrapper = NULL;
    taosGetQitem(tmq->qall, (void**)&pRspWrapper);
1759

1760
    if (pRspWrapper == NULL) {
L
Liu Jicong 已提交
1761
      taosReadAllQitems(tmq->mqueue, tmq->qall);
1762 1763
      taosGetQitem(tmq->qall, (void**)&pRspWrapper);
      if (pRspWrapper == NULL) {
L
Liu Jicong 已提交
1764 1765
        return NULL;
      }
X
Xiaoyu Wang 已提交
1766 1767
    }

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

1770 1771
    if (pRspWrapper->tmqRspType == TMQ_MSG_TYPE__END_RSP) {
      taosFreeQitem(pRspWrapper);
L
Liu Jicong 已提交
1772
      terrno = TSDB_CODE_TQ_NO_COMMITTED_OFFSET;
H
Haojun Liao 已提交
1773
      tscError("consumer:0x%" PRIx64 " unexpected rsp from poll, code:%s", tmq->consumerId, tstrerror(terrno));
L
Liu Jicong 已提交
1774
      return NULL;
1775 1776
    } else if (pRspWrapper->tmqRspType == TMQ_MSG_TYPE__POLL_RSP) {
      SMqPollRspWrapper* pollRspWrapper = (SMqPollRspWrapper*)pRspWrapper;
H
Haojun Liao 已提交
1777

X
Xiaoyu Wang 已提交
1778
      int32_t     consumerEpoch = atomic_load_32(&tmq->epoch);
1779 1780 1781
      SMqDataRsp* pDataRsp = &pollRspWrapper->dataRsp;

      if (pDataRsp->head.epoch == consumerEpoch) {
L
Liu Jicong 已提交
1782
        SMqClientVg* pVg = pollRspWrapper->vgHandle;
1783 1784 1785 1786 1787

        // update the epset
        if (pollRspWrapper->pEpset != NULL) {
          SEp* pEp = GET_ACTIVE_EP(pollRspWrapper->pEpset);
          SEp* pOld = GET_ACTIVE_EP(&(pVg->epSet));
X
Xiaoyu Wang 已提交
1788 1789
          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);
1790 1791 1792
          pVg->epSet = *pollRspWrapper->pEpset;
        }

wmmhello's avatar
wmmhello 已提交
1793 1794 1795
        if(pDataRsp->rspOffset.type != 0){    // if offset is validate
          pVg->currentOffset = pDataRsp->rspOffset;          // update the local offset value only for the returned values.
        }
X
Xiaoyu Wang 已提交
1796
        atomic_store_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE);
H
Haojun Liao 已提交
1797

1798 1799 1800
        char buf[80];
        tFormatOffset(buf, 80, &pDataRsp->rspOffset);
        if (pDataRsp->blockNum == 0) {
X
Xiaoyu Wang 已提交
1801 1802 1803
          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);
1804
          pRspWrapper = tmqFreeRspWrapper(pRspWrapper);
1805
          pVg->emptyBlockReceiveTs = taosGetTimestampMs();
L
Liu Jicong 已提交
1806
          taosFreeQitem(pollRspWrapper);
1807
        } else {  // build rsp
X
Xiaoyu Wang 已提交
1808
          int64_t    numOfRows = 0;
1809
          SMqRspObj* pRsp = tmqBuildRspFromWrapper(pollRspWrapper, pVg, &numOfRows);
H
Haojun Liao 已提交
1810
          tmq->totalRows += numOfRows;
1811
          pVg->emptyBlockReceiveTs = 0;
H
Haojun Liao 已提交
1812
          tscDebug("consumer:0x%" PRIx64 " process poll rsp, vgId:%d, offset:%s, blocks:%d, rows:%" PRId64
H
Haojun Liao 已提交
1813
                   " vg total:%" PRId64 " total:%" PRId64 ", reqId:0x%" PRIx64,
H
Haojun Liao 已提交
1814
                   tmq->consumerId, pVg->vgId, buf, pDataRsp->blockNum, numOfRows, pVg->numOfRows, tmq->totalRows,
H
Haojun Liao 已提交
1815
                   pollRspWrapper->reqId);
1816 1817 1818
          taosFreeQitem(pollRspWrapper);
          return pRsp;
        }
X
Xiaoyu Wang 已提交
1819
      } else {
H
Haojun Liao 已提交
1820
        tscDebug("consumer:0x%" PRIx64 " vgId:%d msg discard since epoch mismatch: msg epoch %d, consumer epoch %d",
1821
                 tmq->consumerId, pollRspWrapper->vgId, pDataRsp->head.epoch, consumerEpoch);
1822
        pRspWrapper = tmqFreeRspWrapper(pRspWrapper);
L
Liu Jicong 已提交
1823 1824
        taosFreeQitem(pollRspWrapper);
      }
1825 1826
    } else if (pRspWrapper->tmqRspType == TMQ_MSG_TYPE__POLL_META_RSP) {
      SMqPollRspWrapper* pollRspWrapper = (SMqPollRspWrapper*)pRspWrapper;
L
Liu Jicong 已提交
1827
      int32_t            consumerEpoch = atomic_load_32(&tmq->epoch);
1828 1829 1830

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

L
Liu Jicong 已提交
1831
      if (pollRspWrapper->metaRsp.head.epoch == consumerEpoch) {
L
Liu Jicong 已提交
1832
        SMqClientVg* pVg = pollRspWrapper->vgHandle;
1833 1834 1835
        if(pollRspWrapper->metaRsp.rspOffset.type != 0){    // if offset is validate
          pVg->currentOffset = pollRspWrapper->metaRsp.rspOffset;
        }
L
Liu Jicong 已提交
1836 1837
        atomic_store_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE);
        // build rsp
L
Liu Jicong 已提交
1838
        SMqMetaRspObj* pRsp = tmqBuildMetaRspFromWrapper(pollRspWrapper);
L
Liu Jicong 已提交
1839 1840 1841
        taosFreeQitem(pollRspWrapper);
        return pRsp;
      } else {
H
Haojun Liao 已提交
1842
        tscDebug("consumer:0x%" PRIx64 " vgId:%d msg discard since epoch mismatch: msg epoch %d, consumer epoch %d",
1843
                 tmq->consumerId, pollRspWrapper->vgId, pollRspWrapper->metaRsp.head.epoch, consumerEpoch);
1844
        pRspWrapper = tmqFreeRspWrapper(pRspWrapper);
L
Liu Jicong 已提交
1845
        taosFreeQitem(pollRspWrapper);
X
Xiaoyu Wang 已提交
1846
      }
1847 1848
    } else if (pRspWrapper->tmqRspType == TMQ_MSG_TYPE__TAOSX_RSP) {
      SMqPollRspWrapper* pollRspWrapper = (SMqPollRspWrapper*)pRspWrapper;
X
Xiaoyu Wang 已提交
1849
      int32_t            consumerEpoch = atomic_load_32(&tmq->epoch);
H
Haojun Liao 已提交
1850

L
Liu Jicong 已提交
1851 1852
      if (pollRspWrapper->taosxRsp.head.epoch == consumerEpoch) {
        SMqClientVg* pVg = pollRspWrapper->vgHandle;
1853 1854 1855
        if(pollRspWrapper->taosxRsp.rspOffset.type != 0){    // if offset is validate
          pVg->currentOffset = pollRspWrapper->taosxRsp.rspOffset;
        }
L
Liu Jicong 已提交
1856
        atomic_store_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE);
H
Haojun Liao 已提交
1857

L
Liu Jicong 已提交
1858
        if (pollRspWrapper->taosxRsp.blockNum == 0) {
H
Haojun Liao 已提交
1859 1860
          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 已提交
1861
          pVg->emptyBlockReceiveTs = taosGetTimestampMs();
1862
          pRspWrapper = tmqFreeRspWrapper(pRspWrapper);
H
Haojun Liao 已提交
1863
          taosFreeQitem(pollRspWrapper);
L
Liu Jicong 已提交
1864
          continue;
H
Haojun Liao 已提交
1865
        } else {
X
Xiaoyu Wang 已提交
1866
          pVg->emptyBlockReceiveTs = 0;  // reset the ts
L
Liu Jicong 已提交
1867
        }
wmmhello's avatar
wmmhello 已提交
1868

L
Liu Jicong 已提交
1869
        // build rsp
X
Xiaoyu Wang 已提交
1870
        void*   pRsp = NULL;
1871
        int64_t numOfRows = 0;
L
Liu Jicong 已提交
1872
        if (pollRspWrapper->taosxRsp.createTableNum == 0) {
1873
          pRsp = tmqBuildRspFromWrapper(pollRspWrapper, pVg, &numOfRows);
L
Liu Jicong 已提交
1874
        } else {
wmmhello's avatar
wmmhello 已提交
1875 1876
          pRsp = tmqBuildTaosxRspFromWrapper(pollRspWrapper);
        }
H
Haojun Liao 已提交
1877

1878 1879
        tmq->totalRows += numOfRows;

H
Haojun Liao 已提交
1880 1881
        char buf[80];
        tFormatOffset(buf, 80, &pVg->currentOffset);
H
Haojun Liao 已提交
1882
        tscDebug("consumer:0x%" PRIx64 " process taosx poll rsp, vgId:%d, offset:%s, blocks:%d, rows:%" PRId64
X
Xiaoyu Wang 已提交
1883
                 ", vg total:%" PRId64 " total:%" PRId64 " reqId:0x%" PRIx64,
H
Haojun Liao 已提交
1884
                 tmq->consumerId, pVg->vgId, buf, pollRspWrapper->dataRsp.blockNum, numOfRows, pVg->numOfRows,
H
Haojun Liao 已提交
1885
                 tmq->totalRows, pollRspWrapper->reqId);
H
Haojun Liao 已提交
1886 1887

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

L
Liu Jicong 已提交
1890
      } else {
H
Haojun Liao 已提交
1891
        tscDebug("consumer:0x%" PRIx64 " vgId:%d msg discard since epoch mismatch: msg epoch %d, consumer epoch %d",
1892
                 tmq->consumerId, pollRspWrapper->vgId, pollRspWrapper->taosxRsp.head.epoch, consumerEpoch);
1893
        pRspWrapper = tmqFreeRspWrapper(pRspWrapper);
L
Liu Jicong 已提交
1894 1895
        taosFreeQitem(pollRspWrapper);
      }
X
Xiaoyu Wang 已提交
1896
    } else {
H
Haojun Liao 已提交
1897 1898
      tscDebug("consumer:0x%" PRIx64 " not data msg received", tmq->consumerId);

X
Xiaoyu Wang 已提交
1899
      bool reset = false;
1900 1901
      tmqHandleNoPollRsp(tmq, pRspWrapper, &reset);
      taosFreeQitem(pRspWrapper);
X
Xiaoyu Wang 已提交
1902
      if (pollIfReset && reset) {
1903
        tscDebug("consumer:0x%" PRIx64 ", reset and repoll", tmq->consumerId);
1904
        tmqPollImpl(tmq, timeout);
X
Xiaoyu Wang 已提交
1905 1906 1907 1908 1909
      }
    }
  }
}

1910
TAOS_RES* tmq_consumer_poll(tmq_t* tmq, int64_t timeout) {
L
Liu Jicong 已提交
1911 1912
  void*   rspObj;
  int64_t startTime = taosGetTimestampMs();
L
Liu Jicong 已提交
1913

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

1917 1918 1919
#if 0
  tmqHandleAllDelayedTask(tmq);
  tmqPollImpl(tmq, timeout);
1920
  rspObj = tmqHandleAllRsp(tmq, timeout, false);
L
Liu Jicong 已提交
1921 1922
  if (rspObj) {
    return (TAOS_RES*)rspObj;
L
fix  
Liu Jicong 已提交
1923
  }
1924
#endif
X
Xiaoyu Wang 已提交
1925

1926
  // in no topic status, delayed task also need to be processed
L
Liu Jicong 已提交
1927
  if (atomic_load_8(&tmq->status) == TMQ_CONSUMER_STATUS__INIT) {
1928
    tscDebug("consumer:0x%" PRIx64 " poll return since consumer is init", tmq->consumerId);
1929
    taosMsleep(500);  //     sleep for a while
1930 1931 1932
    return NULL;
  }

wmmhello's avatar
wmmhello 已提交
1933
  while (atomic_load_8(&tmq->status) == TMQ_CONSUMER_STATUS__RECOVER) {
L
Liu Jicong 已提交
1934
    int32_t retryCnt = 0;
1935
    while (TSDB_CODE_MND_CONSUMER_NOT_READY == doAskEp(tmq)) {
H
Haojun Liao 已提交
1936
      if (retryCnt++ > 40) {
L
Liu Jicong 已提交
1937 1938
        return NULL;
      }
1939

H
Haojun Liao 已提交
1940
      tscDebug("consumer:0x%" PRIx64 " not ready, retry:%d/40 in 500ms", tmq->consumerId, retryCnt);
L
Liu Jicong 已提交
1941 1942 1943 1944
      taosMsleep(500);
    }
  }

X
Xiaoyu Wang 已提交
1945
  while (1) {
L
Liu Jicong 已提交
1946
    tmqHandleAllDelayedTask(tmq);
1947

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

1952
    rspObj = tmqHandleAllRsp(tmq, timeout, false);
L
Liu Jicong 已提交
1953
    if (rspObj) {
1954
      tscDebug("consumer:0x%" PRIx64 " return rsp %p", tmq->consumerId, rspObj);
L
Liu Jicong 已提交
1955
      return (TAOS_RES*)rspObj;
L
Liu Jicong 已提交
1956
    } else if (terrno == TSDB_CODE_TQ_NO_COMMITTED_OFFSET) {
1957
      tscDebug("consumer:0x%" PRIx64 " return null since no committed offset", tmq->consumerId);
L
Liu Jicong 已提交
1958
      return NULL;
X
Xiaoyu Wang 已提交
1959
    }
1960

1961
    if (timeout >= 0) {
L
Liu Jicong 已提交
1962
      int64_t currentTime = taosGetTimestampMs();
1963 1964 1965
      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 已提交
1966
                 tmq->consumerId, tmq->epoch, startTime, currentTime);
X
Xiaoyu Wang 已提交
1967 1968
        return NULL;
      }
1969
      tsem_timewait(&tmq->rspSem, (timeout - elapsedTime));
L
Liu Jicong 已提交
1970 1971
    } else {
      // use tsem_timewait instead of tsem_wait to avoid unexpected stuck
L
Liu Jicong 已提交
1972
      tsem_timewait(&tmq->rspSem, 1000);
X
Xiaoyu Wang 已提交
1973 1974 1975 1976
    }
  }
}

1977 1978 1979 1980 1981 1982 1983 1984 1985 1986 1987 1988 1989 1990
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);
1991
    }
1992
  }
1993

1994 1995
  tscDebug("consumer:0x%" PRIx64 " rows dist end", pTmq->consumerId);
}
1996

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

2001 2002 2003 2004 2005 2006
  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;
2007 2008 2009
      }
    }

L
Liu Jicong 已提交
2010
    int32_t     retryCnt = 0;
2011
    tmq_list_t* lst = tmq_list_new();
L
Liu Jicong 已提交
2012
    while (1) {
2013
      int32_t rsp = tmq_subscribe(tmq, lst);
L
Liu Jicong 已提交
2014 2015 2016 2017 2018 2019 2020 2021
      if (rsp != TSDB_CODE_MND_CONSUMER_NOT_READY || retryCnt > 5) {
        break;
      } else {
        retryCnt++;
        taosMsleep(500);
      }
    }

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

2027
  taosRemoveRef(tmqMgmt.rsetId, tmq->refId);
L
Liu Jicong 已提交
2028
  return 0;
2029
}
L
Liu Jicong 已提交
2030

L
Liu Jicong 已提交
2031 2032
const char* tmq_err2str(int32_t err) {
  if (err == 0) {
L
Liu Jicong 已提交
2033
    return "success";
L
Liu Jicong 已提交
2034
  } else if (err == -1) {
L
Liu Jicong 已提交
2035 2036 2037
    return "fail";
  } else {
    return tstrerror(err);
L
Liu Jicong 已提交
2038 2039
  }
}
L
Liu Jicong 已提交
2040

L
Liu Jicong 已提交
2041 2042 2043 2044 2045
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;
2046 2047
  } else if (TD_RES_TMQ_METADATA(res)) {
    return TMQ_RES_METADATA;
L
Liu Jicong 已提交
2048 2049 2050 2051 2052
  } else {
    return TMQ_RES_INVALID;
  }
}

L
Liu Jicong 已提交
2053
const char* tmq_get_topic_name(TAOS_RES* res) {
L
Liu Jicong 已提交
2054 2055
  if (TD_RES_TMQ(res)) {
    SMqRspObj* pRspObj = (SMqRspObj*)res;
L
Liu Jicong 已提交
2056
    return strchr(pRspObj->topic, '.') + 1;
L
Liu Jicong 已提交
2057 2058 2059
  } else if (TD_RES_TMQ_META(res)) {
    SMqMetaRspObj* pMetaRspObj = (SMqMetaRspObj*)res;
    return strchr(pMetaRspObj->topic, '.') + 1;
2060 2061 2062
  } else if (TD_RES_TMQ_METADATA(res)) {
    SMqTaosxRspObj* pRspObj = (SMqTaosxRspObj*)res;
    return strchr(pRspObj->topic, '.') + 1;
L
Liu Jicong 已提交
2063 2064 2065 2066 2067
  } else {
    return NULL;
  }
}

L
Liu Jicong 已提交
2068 2069 2070 2071
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 已提交
2072 2073 2074
  } else if (TD_RES_TMQ_META(res)) {
    SMqMetaRspObj* pMetaRspObj = (SMqMetaRspObj*)res;
    return strchr(pMetaRspObj->db, '.') + 1;
2075 2076 2077
  } else if (TD_RES_TMQ_METADATA(res)) {
    SMqTaosxRspObj* pRspObj = (SMqTaosxRspObj*)res;
    return strchr(pRspObj->db, '.') + 1;
L
Liu Jicong 已提交
2078 2079 2080 2081 2082
  } else {
    return NULL;
  }
}

L
Liu Jicong 已提交
2083 2084 2085 2086
int32_t tmq_get_vgroup_id(TAOS_RES* res) {
  if (TD_RES_TMQ(res)) {
    SMqRspObj* pRspObj = (SMqRspObj*)res;
    return pRspObj->vgId;
L
Liu Jicong 已提交
2087 2088 2089
  } else if (TD_RES_TMQ_META(res)) {
    SMqMetaRspObj* pMetaRspObj = (SMqMetaRspObj*)res;
    return pMetaRspObj->vgId;
2090
  } else if (TD_RES_TMQ_METADATA(res)) {
2091 2092
    SMqTaosxRspObj* pRspObj = (SMqTaosxRspObj*)res;
    return pRspObj->vgId;
L
Liu Jicong 已提交
2093 2094 2095 2096
  } else {
    return -1;
  }
}
L
Liu Jicong 已提交
2097 2098 2099 2100 2101 2102 2103 2104

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;
    }
2105
    return (const char*)taosArrayGetP(pRspObj->rsp.blockTbName, pRspObj->resIter);
2106 2107
  } else if (TD_RES_TMQ_METADATA(res)) {
    SMqTaosxRspObj* pRspObj = (SMqTaosxRspObj*)res;
L
Liu Jicong 已提交
2108 2109 2110
    if (!pRspObj->rsp.withTbName || pRspObj->rsp.blockTbName == NULL || pRspObj->resIter < 0 ||
        pRspObj->resIter >= pRspObj->rsp.blockNum) {
      return NULL;
2111
    }
L
Liu Jicong 已提交
2112 2113
    return (const char*)taosArrayGetP(pRspObj->rsp.blockTbName, pRspObj->resIter);
  }
L
Liu Jicong 已提交
2114 2115
  return NULL;
}
2116

2117 2118 2119 2120 2121 2122
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
    asyncCommitOffset(tmq, pRes, cb, param);
  }
L
Liu Jicong 已提交
2123 2124
}

2125 2126 2127 2128
static void commitCallBackFn(tmq_t *pTmq, int32_t code, void* param) {
  SSyncCommitInfo* pInfo = (SSyncCommitInfo*) param;
  pInfo->code = code;
  tsem_post(&pInfo->sem);
2129
}
2130

2131 2132 2133 2134 2135 2136 2137 2138 2139
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 已提交
2140 2141
  } else {
    asyncCommitOffset(tmq, pRes, commitCallBackFn, pInfo);
2142 2143
  }

2144 2145
  tsem_wait(&pInfo->sem);
  code = pInfo->code;
H
Haojun Liao 已提交
2146 2147

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

X
Xiaoyu Wang 已提交
2150
  tscDebug("consumer:0x%" PRIx64 " sync commit done, code:%s", tmq->consumerId, tstrerror(code));
2151 2152 2153 2154 2155 2156 2157 2158 2159 2160 2161 2162
  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);
2163
    doUpdateLocalEp(pTmq, head->epoch, &rsp);
2164 2165 2166
    tDeleteSMqAskEpRsp(&rsp);
  }

H
Haojun Liao 已提交
2167
  tsem_post(&pInfo->sem);
2168 2169 2170 2171 2172 2173 2174 2175 2176 2177 2178 2179 2180 2181 2182 2183 2184 2185 2186 2187 2188 2189 2190 2191 2192 2193
}

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 已提交
2194
  tsem_init(&pInfo->sem, 0, 0);
2195 2196

  asyncAskEp(pTmq, updateEpCallbackFn, pInfo);
H
Haojun Liao 已提交
2197
  tsem_wait(&pInfo->sem);
2198 2199

  int32_t code = pInfo->code;
H
Haojun Liao 已提交
2200
  tsem_destroy(&pInfo->sem);
2201 2202 2203 2204 2205
  taosMemoryFree(pInfo);
  return code;
}

void asyncAskEp(tmq_t* pTmq, __tmq_askep_fn_t askEpFn, void* param) {
2206
  SMqAskEpReq req = {0};
2207 2208 2209
  req.consumerId = pTmq->consumerId;
  req.epoch = pTmq->epoch;
  strcpy(req.cgroup, pTmq->groupId);
2210 2211 2212

  int32_t tlen = tSerializeSMqAskEpReq(NULL, 0, &req);
  if (tlen < 0) {
2213 2214 2215
    tscError("consumer:0x%" PRIx64 ", tSerializeSMqAskEpReq failed", pTmq->consumerId);
    askEpFn(pTmq, TSDB_CODE_INVALID_PARA, NULL, param);
    return;
2216 2217 2218 2219
  }

  void* pReq = taosMemoryCalloc(1, tlen);
  if (pReq == NULL) {
2220 2221 2222
    tscError("consumer:0x%" PRIx64 ", failed to malloc askEpReq msg, size:%d", pTmq->consumerId, tlen);
    askEpFn(pTmq, TSDB_CODE_OUT_OF_MEMORY, NULL, param);
    return;
2223 2224 2225
  }

  if (tSerializeSMqAskEpReq(pReq, tlen, &req) < 0) {
2226
    tscError("consumer:0x%" PRIx64 ", tSerializeSMqAskEpReq %d failed", pTmq->consumerId, tlen);
2227
    taosMemoryFree(pReq);
2228 2229 2230

    askEpFn(pTmq, TSDB_CODE_INVALID_PARA, NULL, param);
    return;
2231 2232 2233 2234
  }

  SMqAskEpCbParam* pParam = taosMemoryCalloc(1, sizeof(SMqAskEpCbParam));
  if (pParam == NULL) {
2235
    tscError("consumer:0x%" PRIx64 ", failed to malloc subscribe param", pTmq->consumerId);
2236
    taosMemoryFree(pReq);
2237 2238 2239

    askEpFn(pTmq, TSDB_CODE_OUT_OF_MEMORY, NULL, param);
    return;
2240 2241
  }

2242 2243 2244 2245
  pParam->refId = pTmq->refId;
  pParam->epoch = pTmq->epoch;
  pParam->pUserFn = askEpFn;
  pParam->pParam = param;
2246 2247 2248 2249 2250

  SMsgSendInfo* sendInfo = taosMemoryCalloc(1, sizeof(SMsgSendInfo));
  if (sendInfo == NULL) {
    taosMemoryFree(pParam);
    taosMemoryFree(pReq);
2251 2252
    askEpFn(pTmq, TSDB_CODE_OUT_OF_MEMORY, NULL, param);
    return;
2253 2254
  }

X
Xiaoyu Wang 已提交
2255
  sendInfo->msgInfo = (SDataBuf){.pData = pReq, .len = tlen, .handle = NULL};
2256 2257 2258 2259

  sendInfo->requestId = generateRequestId();
  sendInfo->requestObjRefId = 0;
  sendInfo->param = pParam;
2260
  sendInfo->fp = askEpCallbackFn;
2261 2262
  sendInfo->msgType = TDMT_MND_TMQ_ASK_EP;

2263 2264
  SEpSet epSet = getEpSet_s(&pTmq->pTscObj->pAppInfo->mgmtEp);
  tscDebug("consumer:0x%" PRIx64 " ask ep from mnode, reqId:0x%" PRIx64, pTmq->consumerId, sendInfo->requestId);
2265 2266

  int64_t transporterId = 0;
2267
  asyncSendMsgToServer(pTmq->pTscObj->pAppInfo->pTransporter, &epSet, &transporterId, sendInfo);
2268 2269 2270 2271 2272 2273 2274
}

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

int32_t tmqCommitDone(SMqCommitCbParamSet* pParamSet) {
2275 2276 2277
  int64_t refId = pParamSet->refId;

  tmq_t* tmq = taosAcquireRef(tmqMgmt.rsetId, refId);
2278 2279 2280 2281 2282 2283 2284
  if (tmq == NULL) {
    taosMemoryFree(pParamSet);
    terrno = TSDB_CODE_TMQ_CONSUMER_CLOSED;
    return -1;
  }

  // if no more waiting rsp
2285
  pParamSet->callbackFn(tmq, pParamSet->code, pParamSet->userParam);
2286
  taosMemoryFree(pParamSet);
2287 2288

  taosReleaseRef(tmqMgmt.rsetId, refId);
2289
  return 0;
2290 2291
}

2292
void commitRspCountDown(SMqCommitCbParamSet* pParamSet, int64_t consumerId, const char* pTopic, int32_t vgId) {
2293 2294
  int32_t waitingRspNum = atomic_sub_fetch_32(&pParamSet->waitingRspNum, 1);
  if (waitingRspNum == 0) {
H
Haojun Liao 已提交
2295 2296
    tscDebug("consumer:0x%" PRIx64 " topic:%s vgId:%d all commit-rsp received, commit completed", consumerId, pTopic,
             vgId);
2297
    tmqCommitDone(pParamSet);
H
Haojun Liao 已提交
2298 2299 2300
  } else {
    tscDebug("consumer:0x%" PRIx64 " topic:%s vgId:%d commit-rsp received, remain:%d", consumerId, pTopic, vgId,
             waitingRspNum);
2301 2302
  }
}
2303 2304 2305 2306 2307 2308 2309 2310 2311 2312 2313 2314 2315 2316 2317 2318 2319 2320 2321 2322 2323 2324 2325

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