clientTmq.c 72.9 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"

27
#define EMPTY_BLOCK_POLL_IDLE_DURATION  10
28
#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 97 98
  // poll info
  int64_t       pollCnt;
  int64_t       totalRows;
L
Liu Jicong 已提交
99

L
Liu Jicong 已提交
100
  // timer
101 102 103 104
  tmr_h         hbLiveTimer;
  tmr_h         epTimer;
  tmr_h         reportTimer;
  tmr_h         commitTimer;
H
Haojun Liao 已提交
105 106 107 108 109 110 111
  STscObj*      pTscObj;       // connection
  SArray*       clientTopics;  // SArray<SMqClientTopic>
  STaosQueue*   mqueue;        // queue of rsp
  STaosQall*    qall;
  STaosQueue*   delayedTask;   // delayed task queue for heartbeat and auto commit
  TdThreadMutex lock;          // used to protect the operation on each topic, when updating the epsets.
  tsem_t        rspSem;
L
Liu Jicong 已提交
112 113
};

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

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

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

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

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

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

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

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

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

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

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

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

215
static int32_t doAskEp(tmq_t* tmq);
216 217
static int32_t makeTopicVgroupKey(char* dst, const char* topicName, int32_t vg);
static int32_t tmqCommitDone(SMqCommitCbParamSet* pParamSet);
218 219
static int32_t doSendCommitMsg(tmq_t* tmq, SMqClientVg* pVg, const char* pTopicName, SMqCommitCbParamSet* pParamSet,
                               int32_t index, int32_t totalVgroups);
220 221 222
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);
223

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

231
  conf->withTbName = false;
L
Liu Jicong 已提交
232
  conf->autoCommit = true;
233
  conf->autoCommitInterval = DEFAULT_AUTO_COMMIT_INTERVAL;
234
  conf->resetOffset = TMQ_OFFSET__RESET_EARLIEAST;
235
  conf->hbBgEnable = true;
236

237 238 239
  return conf;
}

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

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

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

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

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

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

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

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

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

327
  if (strcasecmp(key, "enable.heartbeat.background") == 0) {
328 329 330 331 332 333 334 335
//    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");
L
Liu Jicong 已提交
336
      return TMQ_CONF_INVALID;
337
//    }
L
Liu Jicong 已提交
338 339
  }

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

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

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

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

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

L
Liu Jicong 已提交
364
  return TMQ_CONF_UNKNOWN;
365 366 367
}

tmq_list_t* tmq_list_new() {
L
Liu Jicong 已提交
368
  return (tmq_list_t*)taosArrayInit(0, sizeof(void*));
369 370
}

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

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

L
Liu Jicong 已提交
387 388 389 390 391 392 393 394 395 396
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;
}

397 398 399 400 401
static SMqClientVg* foundClientVg(SArray* pTopicList, const char* pName, int32_t vgId, int32_t* index, int32_t* numOfVgroups) {
  int32_t numOfTopics = taosArrayGetSize(pTopicList);
  *index = -1;
  *numOfVgroups = 0;

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

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

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

421 422 423
// 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 已提交
424
static int32_t tmqCommitCb(void* param, SDataBuf* pBuf, int32_t code) {
425
  SMqCommitCbParam*    pParam = (SMqCommitCbParam*)param;
426
  SMqCommitCbParamSet* pParamSet = (SMqCommitCbParamSet*)pParam->params;
H
Haojun Liao 已提交
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
//  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);
//
452
//    commitRspCountDown(pParamSet, pParam->pTmq->consumerId, pParam->topicName, pParam->vgId);
453 454 455 456
//    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

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

  pMsgSendInfo->requestId = generateRequestId();
  pMsgSendInfo->requestObjRefId = 0;
  pMsgSendInfo->param = pParam;
L
Liu Jicong 已提交
535
  pMsgSendInfo->paramFreeFp = taosMemoryFree;
536
  pMsgSendInfo->fp = tmqCommitCb;
L
Liu Jicong 已提交
537
  pMsgSendInfo->msgType = TDMT_VND_TMQ_COMMIT_OFFSET;
L
Liu Jicong 已提交
538

L
Liu Jicong 已提交
539 540 541
  atomic_add_fetch_32(&pParamSet->waitingRspNum, 1);
  atomic_add_fetch_32(&pParamSet->totalRspNum, 1);

H
Haojun Liao 已提交
542
  SEp* pEp = GET_ACTIVE_EP(&pVg->epSet);
543 544 545 546 547 548 549 550
  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 已提交
551

L
Liu Jicong 已提交
552 553
  int64_t transporterId = 0;
  asyncSendMsgToServer(tmq->pTscObj->pAppInfo->pTransporter, &pVg->epSet, &transporterId, pMsgSendInfo);
554 555

  return TSDB_CODE_SUCCESS;
L
Liu Jicong 已提交
556 557
}

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

  SMqCommitCbParamSet* pParamSet = taosMemoryCalloc(1, sizeof(SMqCommitCbParamSet));
  if (pParamSet == NULL) {
587 588
    pCommitFp(tmq, TSDB_CODE_OUT_OF_MEMORY, userParam);
    return;
L
Liu Jicong 已提交
589
  }
H
Haojun Liao 已提交
590

591 592
  pParamSet->refId = tmq->refId;
  pParamSet->epoch = tmq->epoch;
593
  pParamSet->callbackFn = pCommitFp;
L
Liu Jicong 已提交
594
  pParamSet->userParam = userParam;
L
Liu Jicong 已提交
595

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

598 599 600 601
  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 已提交
602
    SMqClientTopic* pTopic = taosArrayGet(tmq->clientTopics, i);
603 604
    if (strcmp(pTopic->topicName, pTopicName) == 0) {
      break;
605
    }
606
  }
607

608
  if (i == numOfTopics) {
H
Haojun Liao 已提交
609
    tscWarn("consumer:0x%" PRIx64 " failed to find the specified topic:%s, total topics:%d", tmq->consumerId, pTopicName,
610 611 612 613 614
            numOfTopics);
    taosMemoryFree(pParamSet);
    pCommitFp(tmq, TSDB_CODE_SUCCESS, userParam);
    return;
  }
L
Liu Jicong 已提交
615

616 617 618 619 620 621 622 623
  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 已提交
624
    }
L
Liu Jicong 已提交
625
  }
L
Liu Jicong 已提交
626

627
  if (j == numOfVgroups) {
H
Haojun Liao 已提交
628
    tscWarn("consumer:0x%" PRIx64 " failed to find the specified vgId:%d, total Vgs:%d, topic:%s", tmq->consumerId, vgId,
629
            numOfVgroups, pTopicName);
L
Liu Jicong 已提交
630
    taosMemoryFree(pParamSet);
631 632
    pCommitFp(tmq, TSDB_CODE_SUCCESS, userParam);
    return;
L
Liu Jicong 已提交
633 634
  }

635 636 637
  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 已提交
638

639 640 641 642 643 644 645 646
    // failed to commit, callback user function directly.
    if (code != TSDB_CODE_SUCCESS) {
      taosMemoryFree(pParamSet);
      pCommitFp(tmq, code, userParam);
    }
  } else { // do not perform commit, callback user function directly.
    taosMemoryFree(pParamSet);
    pCommitFp(tmq, code, userParam);
L
Liu Jicong 已提交
647 648 649
  }
}

650
static void asyncCommitAllOffsets(tmq_t* tmq, tmq_commit_cb* pCommitFp, void* userParam) {
651 652
  SMqCommitCbParamSet* pParamSet = taosMemoryCalloc(1, sizeof(SMqCommitCbParamSet));
  if (pParamSet == NULL) {
653 654
    pCommitFp(tmq, TSDB_CODE_OUT_OF_MEMORY, userParam);
    return;
655
  }
656 657 658

  pParamSet->refId = tmq->refId;
  pParamSet->epoch = tmq->epoch;
659
  pParamSet->callbackFn = pCommitFp;
660 661
  pParamSet->userParam = userParam;

662 663 664
  // init as 1 to prevent concurrency issue
  pParamSet->waitingRspNum = 1;

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

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

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

      if (pVg->currentOffset.type > 0 && !tOffsetEqual(&pVg->currentOffset, &pVg->committedOffset)) {
678
        int32_t code = doSendCommitMsg(tmq, pVg, pTopic->topicName, pParamSet, j, numOfVgroups);
679 680 681 682
        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 已提交
683 684
          continue;
        }
H
Haojun Liao 已提交
685 686 687

        // update the offset value.
        pVg->committedOffset = pVg->currentOffset;
688
      } else {
689
        tscDebug("consumer:0x%" PRIx64 " topic:%s vgId:%d, not commit, current:%" PRId64 ", ordinal:%d/%d",
690
                 tmq->consumerId, pTopic->topicName, pVg->vgId, pVg->currentOffset.version, j + 1, numOfVgroups);
691 692 693 694
      }
    }
  }

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

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

L
Liu Jicong 已提交
705
  // count down since waiting rsp num init as 1
706
  commitRspCountDown(pParamSet, tmq->consumerId, "", 0);
707 708
}

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

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

void tmqAssignDelayedCommitTask(void* param, void* tmrId) {
727
  int64_t refId = *(int64_t*)param;
728
  generateTimedTask(refId, TMQ_DELAYED_TASK__COMMIT);
729
  taosMemoryFree(param);
L
Liu Jicong 已提交
730 731 732
}

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

  taosReleaseRef(tmqMgmt.rsetId, refId);
743
  taosMemoryFree(param);
L
Liu Jicong 已提交
744 745
}

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

void tmqSendHbReq(void* param, void* tmrId) {
755
  int64_t refId = *(int64_t*)param;
756

757 758
  tmq_t*  tmq = taosAcquireRef(tmqMgmt.rsetId, refId);
  if (tmq == NULL) {
L
Liu Jicong 已提交
759
    taosMemoryFree(param);
760 761
    return;
  }
D
dapan1121 已提交
762 763 764 765 766

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

L
Liu Jicong 已提交
767
  int32_t tlen = tSerializeSMqHbReq(NULL, 0, &req);
D
dapan1121 已提交
768 769
  if (tlen < 0) {
    tscError("tSerializeSMqHbReq failed");
770
    goto OVER;
D
dapan1121 已提交
771
  }
772

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

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

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

791 792
  sendInfo->msgInfo = (SDataBuf){
      .pData = pReq,
D
dapan1121 已提交
793
      .len = tlen,
794 795 796 797 798 799 800
      .handle = NULL,
  };

  sendInfo->requestId = generateRequestId();
  sendInfo->requestObjRefId = 0;
  sendInfo->param = NULL;
  sendInfo->fp = tmqHbCb;
L
Liu Jicong 已提交
801
  sendInfo->msgType = TDMT_MND_TMQ_HB;
802 803 804 805 806 807 808

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

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

OVER:
809
  taosTmrReset(tmqSendHbReq, 1000, param, tmqMgmt.timer, &tmq->hbLiveTimer);
810
  taosReleaseRef(tmqMgmt.rsetId, refId);
811 812
}

813 814 815 816 817 818
static void defaultCommitCbFn(tmq_t* pTmq, int32_t code, void* param) {
  if (code != 0) {
    tscDebug("consumer:0x%"PRIx64", failed to commit offset, code:%s", pTmq->consumerId, tstrerror(code));
  }
}

819
int32_t tmqHandleAllDelayedTask(tmq_t* pTmq) {
L
Liu Jicong 已提交
820
  STaosQall* qall = taosAllocateQall();
821
  taosReadAllQitems(pTmq->delayedTask, qall);
L
Liu Jicong 已提交
822

823 824 825 826
  if (qall->numOfItems == 0) {
    taosFreeQall(qall);
    return TSDB_CODE_SUCCESS;
  }
827

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

832
  while (pTaskType != NULL) {
833
    if (*pTaskType == TMQ_DELAYED_TASK__ASK_EP) {
834
      asyncAskEp(pTmq, addToQueueCallbackFn, NULL);
835 836

      int64_t* pRefId = taosMemoryMalloc(sizeof(int64_t));
837
      *pRefId = pTmq->refId;
838

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

      asyncCommitAllOffsets(pTmq, pCallbackFn, pTmq->commitCbUserParam);
845
      int64_t* pRefId = taosMemoryMalloc(sizeof(int64_t));
846
      *pRefId = pTmq->refId;
847

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

L
Liu Jicong 已提交
854
    taosFreeQitem(pTaskType);
855
    taosGetQitem(qall, (void**)&pTaskType);
L
Liu Jicong 已提交
856
  }
857

L
Liu Jicong 已提交
858 859 860 861
  taosFreeQall(qall);
  return 0;
}

862
static void* tmqFreeRspWrapper(SMqRspWrapper* rspWrapper) {
L
Liu Jicong 已提交
863 864 865 866 867 868 869
  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;
870 871
    taosMemoryFreeClear(pRsp->pEpset);

L
Liu Jicong 已提交
872 873 874 875 876 877
    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;
878 879
    taosMemoryFreeClear(pRsp->pEpset);

L
Liu Jicong 已提交
880 881 882
    taosMemoryFree(pRsp->metaRsp.metaRsp);
  } else if (rspWrapper->tmqRspType == TMQ_MSG_TYPE__TAOSX_RSP) {
    SMqPollRspWrapper* pRsp = (SMqPollRspWrapper*)rspWrapper;
883 884
    taosMemoryFreeClear(pRsp->pEpset);

L
Liu Jicong 已提交
885 886 887 888 889 890 891 892
    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);
  }
893 894

  return NULL;
L
Liu Jicong 已提交
895 896
}

L
Liu Jicong 已提交
897
void tmqClearUnhandleMsg(tmq_t* tmq) {
L
Liu Jicong 已提交
898
  SMqRspWrapper* rspWrapper = NULL;
L
Liu Jicong 已提交
899
  while (1) {
L
Liu Jicong 已提交
900 901 902 903 904
    taosGetQitem(tmq->qall, (void**)&rspWrapper);
    if (rspWrapper) {
      tmqFreeRspWrapper(rspWrapper);
      taosFreeQitem(rspWrapper);
    } else {
L
Liu Jicong 已提交
905
      break;
L
Liu Jicong 已提交
906
    }
L
Liu Jicong 已提交
907 908
  }

L
Liu Jicong 已提交
909
  rspWrapper = NULL;
L
Liu Jicong 已提交
910 911
  taosReadAllQitems(tmq->mqueue, tmq->qall);
  while (1) {
L
Liu Jicong 已提交
912 913 914 915 916
    taosGetQitem(tmq->qall, (void**)&rspWrapper);
    if (rspWrapper) {
      tmqFreeRspWrapper(rspWrapper);
      taosFreeQitem(rspWrapper);
    } else {
L
Liu Jicong 已提交
917
      break;
L
Liu Jicong 已提交
918
    }
L
Liu Jicong 已提交
919 920 921
  }
}

D
dapan1121 已提交
922
int32_t tmqSubscribeCb(void* param, SDataBuf* pMsg, int32_t code) {
L
Liu Jicong 已提交
923 924
  SMqSubscribeCbParam* pParam = (SMqSubscribeCbParam*)param;
  pParam->rspErr = code;
dengyihao's avatar
dengyihao 已提交
925 926

  taosMemoryFree(pMsg->pEpSet);
L
Liu Jicong 已提交
927 928 929
  tsem_post(&pParam->rspSem);
  return 0;
}
930

L
Liu Jicong 已提交
931
int32_t tmq_subscription(tmq_t* tmq, tmq_list_t** topics) {
X
Xiaoyu Wang 已提交
932 933 934 935
  if (*topics == NULL) {
    *topics = tmq_list_new();
  }
  for (int i = 0; i < taosArrayGetSize(tmq->clientTopics); i++) {
L
Liu Jicong 已提交
936
    SMqClientTopic* topic = taosArrayGet(tmq->clientTopics, i);
L
Liu Jicong 已提交
937
    tmq_list_append(*topics, strchr(topic->topicName, '.') + 1);
X
Xiaoyu Wang 已提交
938
  }
L
Liu Jicong 已提交
939
  return 0;
X
Xiaoyu Wang 已提交
940 941
}

L
Liu Jicong 已提交
942
int32_t tmq_unsubscribe(tmq_t* tmq) {
L
Liu Jicong 已提交
943 944
  int32_t     rsp;
  int32_t     retryCnt = 0;
L
Liu Jicong 已提交
945
  tmq_list_t* lst = tmq_list_new();
L
Liu Jicong 已提交
946 947 948 949 950 951 952 953 954 955
  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 已提交
956 957
  tmq_list_destroy(lst);
  return rsp;
X
Xiaoyu Wang 已提交
958 959
}

960 961 962 963 964 965
static void freeClientVgImpl(void* param) {
  SMqClientTopic* pTopic = param;
  taosMemoryFreeClear(pTopic->schema.pSchema);
  taosArrayDestroy(pTopic->vgs);
}

966
void tmqFreeImpl(void* handle) {
967 968
  tmq_t*  tmq = (tmq_t*)handle;
  int64_t id = tmq->consumerId;
L
Liu Jicong 已提交
969

970
  // TODO stop timer
L
Liu Jicong 已提交
971 972 973 974
  if (tmq->mqueue) {
    tmqClearUnhandleMsg(tmq);
    taosCloseQueue(tmq->mqueue);
  }
L
Liu Jicong 已提交
975

H
Haojun Liao 已提交
976 977 978 979 980
  if (tmq->delayedTask) {
    taosCloseQueue(tmq->delayedTask);
  }

  taosFreeQall(tmq->qall);
981
  tsem_destroy(&tmq->rspSem);
H
Haojun Liao 已提交
982
  taosThreadMutexDestroy(&tmq->lock);
L
Liu Jicong 已提交
983

984
  taosArrayDestroyEx(tmq->clientTopics, freeClientVgImpl);
985 986
  taos_close_internal(tmq->pTscObj);
  taosMemoryFree(tmq);
987 988

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

991 992 993 994 995 996 997 998 999
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);
1000
  if (tmqMgmt.rsetId < 0) {
1001 1002 1003 1004
    tmqInitRes = terrno;
  }
}

L
Liu Jicong 已提交
1005
tmq_t* tmq_consumer_new(tmq_conf_t* conf, char* errstr, int32_t errstrLen) {
1006 1007 1008 1009
  taosThreadOnce(&tmqInit, tmqMgmtInit);
  if (tmqInitRes != 0) {
    terrno = tmqInitRes;
    return NULL;
L
Liu Jicong 已提交
1010 1011
  }

L
Liu Jicong 已提交
1012 1013
  tmq_t* pTmq = taosMemoryCalloc(1, sizeof(tmq_t));
  if (pTmq == NULL) {
L
Liu Jicong 已提交
1014
    terrno = TSDB_CODE_OUT_OF_MEMORY;
1015
    tscError("failed to create consumer, groupId:%s, code:%s", conf->groupId, terrstr());
L
Liu Jicong 已提交
1016 1017
    return NULL;
  }
L
Liu Jicong 已提交
1018

L
Liu Jicong 已提交
1019 1020 1021
  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 已提交
1022 1023 1024
  pTmq->clientTopics = taosArrayInit(0, sizeof(SMqClientTopic));
  pTmq->mqueue = taosOpenQueue();
  pTmq->delayedTask = taosOpenQueue();
H
Haojun Liao 已提交
1025
  pTmq->qall = taosAllocateQall();
L
Liu Jicong 已提交
1026

H
Haojun Liao 已提交
1027
  taosThreadMutexInit(&pTmq->lock, NULL);
X
Xiaoyu Wang 已提交
1028 1029
  if (pTmq->clientTopics == NULL || pTmq->mqueue == NULL || pTmq->qall == NULL || pTmq->delayedTask == NULL ||
      conf->groupId[0] == 0) {
L
Liu Jicong 已提交
1030
    terrno = TSDB_CODE_OUT_OF_MEMORY;
1031
    tscError("consumer:0x%" PRIx64 " setup failed since %s, groupId:%s", pTmq->consumerId, terrstr(),
S
Shengliang Guan 已提交
1032
             pTmq->groupId);
1033
    goto _failed;
L
Liu Jicong 已提交
1034
  }
L
Liu Jicong 已提交
1035

L
Liu Jicong 已提交
1036 1037
  // init status
  pTmq->status = TMQ_CONSUMER_STATUS__INIT;
L
Liu Jicong 已提交
1038 1039
  pTmq->pollCnt = 0;
  pTmq->epoch = 0;
L
Liu Jicong 已提交
1040

L
Liu Jicong 已提交
1041 1042 1043
  // set conf
  strcpy(pTmq->clientId, conf->clientId);
  strcpy(pTmq->groupId, conf->groupId);
1044
  pTmq->withTbName = conf->withTbName;
L
Liu Jicong 已提交
1045
  pTmq->useSnapshot = conf->snapEnable;
L
Liu Jicong 已提交
1046
  pTmq->autoCommit = conf->autoCommit;
L
Liu Jicong 已提交
1047
  pTmq->autoCommitInterval = conf->autoCommitInterval;
L
Liu Jicong 已提交
1048 1049
  pTmq->commitCb = conf->commitCb;
  pTmq->commitCbUserParam = conf->commitCbUserParam;
L
Liu Jicong 已提交
1050 1051
  pTmq->resetOffsetCfg = conf->resetOffset;

1052 1053
  pTmq->hbBgEnable = conf->hbBgEnable;

L
Liu Jicong 已提交
1054
  // assign consumerId
L
Liu Jicong 已提交
1055
  pTmq->consumerId = tGenIdPI64();
X
Xiaoyu Wang 已提交
1056

L
Liu Jicong 已提交
1057 1058
  // init semaphore
  if (tsem_init(&pTmq->rspSem, 0, 0) != 0) {
1059
    tscError("consumer:0x %" PRIx64 " setup failed since %s, consumer group %s", pTmq->consumerId, terrstr(),
S
Shengliang Guan 已提交
1060
             pTmq->groupId);
1061
    goto _failed;
L
Liu Jicong 已提交
1062
  }
L
Liu Jicong 已提交
1063

L
Liu Jicong 已提交
1064 1065 1066
  // init connection
  pTmq->pTscObj = taos_connect_internal(conf->ip, user, pass, NULL, NULL, conf->port, CONN_TYPE__TMQ);
  if (pTmq->pTscObj == NULL) {
1067
    tscError("consumer:0x%" PRIx64 " setup failed since %s, groupId:%s", pTmq->consumerId, terrstr(), pTmq->groupId);
L
Liu Jicong 已提交
1068
    tsem_destroy(&pTmq->rspSem);
1069
    goto _failed;
L
Liu Jicong 已提交
1070
  }
L
Liu Jicong 已提交
1071

1072 1073
  pTmq->refId = taosAddRef(tmqMgmt.rsetId, pTmq);
  if (pTmq->refId < 0) {
1074
    goto _failed;
1075 1076
  }

1077
  if (pTmq->hbBgEnable) {
L
Liu Jicong 已提交
1078 1079
    int64_t* pRefId = taosMemoryMalloc(sizeof(int64_t));
    *pRefId = pTmq->refId;
1080
    pTmq->hbLiveTimer = taosTmrStart(tmqSendHbReq, 1000, pRefId, tmqMgmt.timer);
1081 1082
  }

1083 1084 1085
  char         buf[80] = {0};
  STqOffsetVal offset = {.type = pTmq->resetOffsetCfg};
  tFormatOffset(buf, tListLen(buf), &offset);
1086 1087
  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,
1088
          pTmq->hbBgEnable);
L
Liu Jicong 已提交
1089

1090
  return pTmq;
1091

1092 1093
_failed:
  tmqFreeImpl(pTmq);
L
Liu Jicong 已提交
1094
  return NULL;
1095 1096
}

L
Liu Jicong 已提交
1097
int32_t tmq_subscribe(tmq_t* tmq, const tmq_list_t* topic_list) {
1098
  const int32_t   MAX_RETRY_COUNT = 120 * 2;  // let's wait for 2 mins at most
L
Liu Jicong 已提交
1099 1100 1101
  const SArray*   container = &topic_list->container;
  int32_t         sz = taosArrayGetSize(container);
  void*           buf = NULL;
L
Liu Jicong 已提交
1102
  SMsgSendInfo*   sendInfo = NULL;
L
Liu Jicong 已提交
1103
  SCMSubscribeReq req = {0};
1104
  int32_t         code = 0;
1105

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

1108
  req.consumerId = tmq->consumerId;
L
Liu Jicong 已提交
1109
  tstrncpy(req.clientId, tmq->clientId, 256);
L
Liu Jicong 已提交
1110
  tstrncpy(req.cgroup, tmq->groupId, TSDB_CGROUP_LEN);
1111 1112
  req.topicNames = taosArrayInit(sz, sizeof(void*));

1113 1114 1115 1116
  if (req.topicNames == NULL) {
    code = TSDB_CODE_OUT_OF_MEMORY;
    goto FAIL;
  }
1117

L
Liu Jicong 已提交
1118 1119
  for (int32_t i = 0; i < sz; i++) {
    char* topic = taosArrayGetP(container, i);
1120 1121

    SName name = {0};
L
Liu Jicong 已提交
1122 1123 1124 1125
    tNameSetDbName(&name, tmq->pTscObj->acctId, topic, strlen(topic));
    char* topicFName = taosMemoryCalloc(1, TSDB_TOPIC_FNAME_LEN);
    if (topicFName == NULL) {
      goto FAIL;
1126 1127
    }

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

    taosArrayPush(req.topicNames, &topicFName);
1132 1133
  }

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

L
Liu Jicong 已提交
1136
  buf = taosMemoryMalloc(tlen);
1137 1138 1139 1140
  if (buf == NULL) {
    code = TSDB_CODE_OUT_OF_MEMORY;
    goto FAIL;
  }
L
Liu Jicong 已提交
1141

1142 1143 1144
  void* abuf = buf;
  tSerializeSCMSubscribeReq(&abuf, &req);

L
Liu Jicong 已提交
1145
  sendInfo = taosMemoryCalloc(1, sizeof(SMsgSendInfo));
1146 1147 1148 1149
  if (sendInfo == NULL) {
    code = TSDB_CODE_OUT_OF_MEMORY;
    goto FAIL;
  }
1150

X
Xiaoyu Wang 已提交
1151
  SMqSubscribeCbParam param = {
L
Liu Jicong 已提交
1152
      .rspErr = 0,
1153 1154
      .refId = tmq->refId,
      .epoch = tmq->epoch,
X
Xiaoyu Wang 已提交
1155
  };
L
Liu Jicong 已提交
1156

1157 1158 1159
  if (tsem_init(&param.rspSem, 0, 0) != 0) {
    goto FAIL;
  }
L
Liu Jicong 已提交
1160 1161

  sendInfo->msgInfo = (SDataBuf){
X
Xiaoyu Wang 已提交
1162 1163 1164 1165
      .pData = buf,
      .len = tlen,
      .handle = NULL,
  };
1166

L
Liu Jicong 已提交
1167 1168
  sendInfo->requestId = generateRequestId();
  sendInfo->requestObjRefId = 0;
L
Liu Jicong 已提交
1169 1170
  sendInfo->param = &param;
  sendInfo->fp = tmqSubscribeCb;
L
Liu Jicong 已提交
1171
  sendInfo->msgType = TDMT_MND_TMQ_SUBSCRIBE;
L
Liu Jicong 已提交
1172

1173 1174 1175 1176 1177
  SEpSet epSet = getEpSet_s(&tmq->pTscObj->pAppInfo->mgmtEp);

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

L
Liu Jicong 已提交
1178 1179
  // avoid double free if msg is sent
  buf = NULL;
L
Liu Jicong 已提交
1180
  sendInfo = NULL;
L
Liu Jicong 已提交
1181

L
Liu Jicong 已提交
1182 1183
  tsem_wait(&param.rspSem);
  tsem_destroy(&param.rspSem);
1184

1185 1186 1187 1188
  if (param.rspErr != 0) {
    code = param.rspErr;
    goto FAIL;
  }
L
Liu Jicong 已提交
1189

L
Liu Jicong 已提交
1190
  int32_t retryCnt = 0;
1191
  while (TSDB_CODE_MND_CONSUMER_NOT_READY == doAskEp(tmq)) {
1192
    if (retryCnt++ > MAX_RETRY_COUNT) {
L
Liu Jicong 已提交
1193 1194
      goto FAIL;
    }
1195

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

1200 1201
  // init ep timer
  if (tmq->epTimer == NULL) {
1202 1203 1204
    int64_t* pRefId1 = taosMemoryMalloc(sizeof(int64_t));
    *pRefId1 = tmq->refId;
    tmq->epTimer = taosTmrStart(tmqAssignAskEpTask, 1000, pRefId1, tmqMgmt.timer);
1205
  }
L
Liu Jicong 已提交
1206 1207

  // init auto commit timer
1208
  if (tmq->autoCommit && tmq->commitTimer == NULL) {
1209 1210 1211
    int64_t* pRefId2 = taosMemoryMalloc(sizeof(int64_t));
    *pRefId2 = tmq->refId;
    tmq->commitTimer = taosTmrStart(tmqAssignDelayedCommitTask, tmq->autoCommitInterval, pRefId2, tmqMgmt.timer);
L
Liu Jicong 已提交
1212 1213
  }

L
Liu Jicong 已提交
1214
FAIL:
L
Liu Jicong 已提交
1215
  taosArrayDestroyP(req.topicNames, taosMemoryFree);
L
Liu Jicong 已提交
1216
  taosMemoryFree(buf);
L
Liu Jicong 已提交
1217
  taosMemoryFree(sendInfo);
L
Liu Jicong 已提交
1218

L
Liu Jicong 已提交
1219
  return code;
1220 1221
}

L
Liu Jicong 已提交
1222
void tmq_conf_set_auto_commit_cb(tmq_conf_t* conf, tmq_commit_cb* cb, void* param) {
1223
  conf->commitCb = cb;
L
Liu Jicong 已提交
1224
  conf->commitCbUserParam = param;
L
Liu Jicong 已提交
1225
}
1226

D
dapan1121 已提交
1227
int32_t tmqPollCb(void* param, SDataBuf* pMsg, int32_t code) {
X
Xiaoyu Wang 已提交
1228
  SMqPollCbParam* pParam = (SMqPollCbParam*)param;
1229 1230

  int64_t         refId = pParam->refId;
X
Xiaoyu Wang 已提交
1231
  SMqClientVg*    pVg = pParam->pVg;
L
Liu Jicong 已提交
1232
  SMqClientTopic* pTopic = pParam->pTopic;
1233

1234
  tmq_t* tmq = taosAcquireRef(tmqMgmt.rsetId, refId);
1235 1236 1237 1238
  if (tmq == NULL) {
    tsem_destroy(&pParam->rspSem);
    taosMemoryFree(pParam);
    taosMemoryFree(pMsg->pData);
dengyihao's avatar
dengyihao 已提交
1239
    taosMemoryFree(pMsg->pEpSet);
1240 1241 1242 1243
    terrno = TSDB_CODE_TMQ_CONSUMER_CLOSED;
    return -1;
  }

H
Haojun Liao 已提交
1244 1245 1246 1247
  int32_t  epoch = pParam->epoch;
  int32_t  vgId = pParam->vgId;
  uint64_t requestId = pParam->requestId;

L
Liu Jicong 已提交
1248
  taosMemoryFree(pParam);
H
Haojun Liao 已提交
1249

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

L
Liu Jicong 已提交
1254
    if (pMsg->pData) taosMemoryFree(pMsg->pData);
dengyihao's avatar
dengyihao 已提交
1255 1256
    if (pMsg->pEpSet) taosMemoryFree(pMsg->pEpSet);

H
Haojun Liao 已提交
1257
    // in case of consumer mismatch, wait for 500ms and retry
L
Liu Jicong 已提交
1258
    if (code == TSDB_CODE_TMQ_CONSUMER_MISMATCH) {
1259
      taosMsleep(500);
L
Liu Jicong 已提交
1260
      atomic_store_8(&tmq->status, TMQ_CONSUMER_STATUS__RECOVER);
1261
      tscDebug("consumer:0x%" PRIx64" wait for the re-balance, wait for 500ms and set status to be RECOVER", tmq->consumerId);
H
Haojun Liao 已提交
1262
    } else if (code == TSDB_CODE_TQ_NO_COMMITTED_OFFSET) {
S
Shengliang Guan 已提交
1263
      SMqPollRspWrapper* pRspWrapper = taosAllocateQitem(sizeof(SMqPollRspWrapper), DEF_QITEM, 0);
L
Liu Jicong 已提交
1264
      if (pRspWrapper == NULL) {
H
Haojun Liao 已提交
1265 1266
        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 已提交
1267 1268
        goto CREATE_MSG_FAIL;
      }
H
Haojun Liao 已提交
1269

L
Liu Jicong 已提交
1270 1271
      pRspWrapper->tmqRspType = TMQ_MSG_TYPE__END_RSP;
      taosWriteQitem(tmq->mqueue, pRspWrapper);
1272 1273
    }else if(code == TSDB_CODE_WAL_LOG_NOT_EXIST){    //poll data while insert
      taosMsleep(500);
L
Liu Jicong 已提交
1274
    }
H
Haojun Liao 已提交
1275

L
fix txn  
Liu Jicong 已提交
1276
    goto CREATE_MSG_FAIL;
L
Liu Jicong 已提交
1277 1278
  }

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

1286
    tsem_post(&tmq->rspSem);
1287 1288
    taosReleaseRef(tmqMgmt.rsetId, refId);

L
Liu Jicong 已提交
1289
    taosMemoryFree(pMsg->pData);
dengyihao's avatar
dengyihao 已提交
1290
    taosMemoryFree(pMsg->pEpSet);
X
Xiaoyu Wang 已提交
1291 1292 1293 1294
    return 0;
  }

  if (msgEpoch != tmqEpoch) {
H
Haojun Liao 已提交
1295 1296
    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 已提交
1297 1298
  }

L
Liu Jicong 已提交
1299 1300 1301
  // handle meta rsp
  int8_t rspType = ((SMqRspHead*)pMsg->pData)->mqMsgType;

S
Shengliang Guan 已提交
1302
  SMqPollRspWrapper* pRspWrapper = taosAllocateQitem(sizeof(SMqPollRspWrapper), DEF_QITEM, 0);
L
Liu Jicong 已提交
1303
  if (pRspWrapper == NULL) {
L
Liu Jicong 已提交
1304
    taosMemoryFree(pMsg->pData);
dengyihao's avatar
dengyihao 已提交
1305
    taosMemoryFree(pMsg->pEpSet);
H
Haojun Liao 已提交
1306
    tscWarn("consumer:0x%"PRIx64" msg discard from vgId:%d, epoch %d since out of memory", tmq->consumerId, vgId, epoch);
L
fix txn  
Liu Jicong 已提交
1307
    goto CREATE_MSG_FAIL;
L
Liu Jicong 已提交
1308
  }
L
Liu Jicong 已提交
1309

L
Liu Jicong 已提交
1310
  pRspWrapper->tmqRspType = rspType;
L
Liu Jicong 已提交
1311 1312
  pRspWrapper->vgHandle = pVg;
  pRspWrapper->topicHandle = pTopic;
H
Haojun Liao 已提交
1313
  pRspWrapper->reqId = requestId;
1314
  pRspWrapper->pEpset = pMsg->pEpSet;
1315
  pRspWrapper->vgId = pVg->vgId;
L
Liu Jicong 已提交
1316

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

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

L
Liu Jicong 已提交
1345
  taosMemoryFree(pMsg->pData);
H
Haojun Liao 已提交
1346
  taosWriteQitem(tmq->mqueue, pRspWrapper);
L
Liu Jicong 已提交
1347

H
Haojun Liao 已提交
1348 1349 1350
  tscDebug("consumer:0x%" PRIx64 " put poll res into mqueue, type:%d, vgId:%d, total in queue:%d, reqId:0x%" PRIx64,
           tmq->consumerId, rspType, vgId, tmq->mqueue->numOfItems, requestId);

1351
  tsem_post(&tmq->rspSem);
1352 1353
  taosReleaseRef(tmqMgmt.rsetId, refId);

L
Liu Jicong 已提交
1354
  return 0;
H
Haojun Liao 已提交
1355

L
fix txn  
Liu Jicong 已提交
1356
CREATE_MSG_FAIL:
L
Liu Jicong 已提交
1357
  if (epoch == tmq->epoch) {
L
Liu Jicong 已提交
1358 1359
    atomic_store_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE);
  }
H
Haojun Liao 已提交
1360

1361
  tsem_post(&tmq->rspSem);
1362 1363
  taosReleaseRef(tmqMgmt.rsetId, refId);

L
Liu Jicong 已提交
1364
  return -1;
1365 1366
}

H
Haojun Liao 已提交
1367 1368 1369 1370 1371
typedef struct SVgroupSaveInfo {
  STqOffsetVal offset;
  int64_t      numOfRows;
} SVgroupSaveInfo;

H
Haojun Liao 已提交
1372 1373 1374 1375 1376 1377 1378 1379 1380 1381 1382 1383 1384 1385 1386 1387 1388
static void initClientTopicFromRsp(SMqClientTopic* pTopic, SMqSubTopicEp* pTopicEp, SHashObj* pVgOffsetHashMap,
                                   tmq_t* tmq) {
  pTopic->schema = pTopicEp->schema;
  pTopicEp->schema.nCols = 0;
  pTopicEp->schema.pSchema = NULL;

  char vgKey[TSDB_TOPIC_FNAME_LEN + 22];
  int32_t vgNumGet = taosArrayGetSize(pTopicEp->vgs);

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

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

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

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

H
Haojun Liao 已提交
1393
    int64_t numOfRows = 0;
H
Haojun Liao 已提交
1394
    STqOffsetVal  offsetNew = {.type = tmq->resetOffsetCfg};
H
Haojun Liao 已提交
1395 1396 1397
    if (pInfo != NULL) {
      offsetNew = pInfo->offset;
      numOfRows = pInfo->numOfRows;
H
Haojun Liao 已提交
1398 1399 1400 1401 1402 1403 1404 1405 1406
    }

    SMqClientVg clientVg = {
        .pollCnt = 0,
        .currentOffset = offsetNew,
        .vgId = pVgEp->vgId,
        .epSet = pVgEp->epSet,
        .vgStatus = TMQ_VG_STATUS__IDLE,
        .vgSkipCnt = 0,
H
Haojun Liao 已提交
1407
        .emptyBlockReceiveTs = 0,
H
Haojun Liao 已提交
1408
        .numOfRows = numOfRows,
H
Haojun Liao 已提交
1409 1410 1411 1412 1413 1414 1415 1416 1417 1418 1419 1420 1421 1422 1423 1424
    };

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

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

  taosArrayDestroy(pTopic->vgs);
}

static bool tmqUpdateEp(tmq_t* tmq, int32_t epoch, const SMqAskEpRsp* pRsp) {
1425 1426
  bool set = false;

1427
  int32_t topicNumCur = taosArrayGetSize(tmq->clientTopics);
1428
  int32_t topicNumGet = taosArrayGetSize(pRsp->topics);
1429

X
Xiaoyu Wang 已提交
1430 1431
  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",
1432
           tmq->consumerId, tmq->epoch, epoch, topicNumGet, topicNumCur);
1433 1434 1435 1436 1437 1438

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

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

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

L
Liu Jicong 已提交
1456
        char buf[80];
L
Liu Jicong 已提交
1457
        tFormatOffset(buf, 80, &pVgCur->currentOffset);
H
Haojun Liao 已提交
1458
        tscDebug("consumer:0x%" PRIx64 ", epoch:%d vgId:%d vgKey:%s, offset:%s", tmq->consumerId, epoch,
L
Liu Jicong 已提交
1459
                 pVgCur->vgId, vgKey, buf);
H
Haojun Liao 已提交
1460 1461 1462

        SVgroupSaveInfo info = {.offset = pVgCur->currentOffset, .numOfRows = pVgCur->numOfRows};
        taosHashPut(pVgOffsetHashMap, vgKey, strlen(vgKey), &info, sizeof(SVgroupSaveInfo));
1463 1464 1465 1466 1467 1468 1469
      }
    }
  }

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

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

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

H
Haojun Liao 已提交
1482 1483
  tmq->clientTopics = newTopics;
  taosThreadMutexUnlock(&tmq->lock);
1484

H
Haojun Liao 已提交
1485 1486
  int8_t flag = (topicNumGet == 0)? TMQ_CONSUMER_STATUS__NO_TOPIC:TMQ_CONSUMER_STATUS__READY;
  atomic_store_8(&tmq->status, flag);
X
Xiaoyu Wang 已提交
1487
  atomic_store_32(&tmq->epoch, epoch);
H
Haojun Liao 已提交
1488

1489
  tscDebug("consumer:0x%" PRIx64 " update topic info completed", tmq->consumerId);
X
Xiaoyu Wang 已提交
1490 1491 1492
  return set;
}

1493
int32_t askEpCallbackFn(void* param, SDataBuf* pMsg, int32_t code) {
1494
  SMqAskEpCbParam* pParam = (SMqAskEpCbParam*)param;
1495
  tmq_t* tmq = taosAcquireRef(tmqMgmt.rsetId, pParam->refId);
1496 1497

  if (tmq == NULL) {
1498 1499 1500
    terrno = TSDB_CODE_TMQ_CONSUMER_CLOSED;
    pParam->pUserFn(tmq, terrno, NULL, pParam->pParam);

1501
    taosMemoryFree(pMsg->pData);
dengyihao's avatar
dengyihao 已提交
1502
    taosMemoryFree(pMsg->pEpSet);
1503 1504
    taosMemoryFree(pParam);
    return terrno;
1505 1506
  }

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

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

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

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

dengyihao's avatar
dengyihao 已提交
1544
  taosMemoryFree(pMsg->pEpSet);
L
Liu Jicong 已提交
1545
  taosMemoryFree(pMsg->pData);
1546
  taosMemoryFree(pParam);
L
Liu Jicong 已提交
1547
  return code;
1548 1549
}

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

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

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

1578
SMqRspObj* tmqBuildRspFromWrapper(SMqPollRspWrapper* pWrapper, SMqClientVg* pVg, int64_t* numOfRows) {
L
Liu Jicong 已提交
1579 1580
  SMqRspObj* pRspObj = taosMemoryCalloc(1, sizeof(SMqRspObj));
  pRspObj->resType = RES_TYPE__TMQ;
1581

1582
  (*numOfRows) = 0;
L
Liu Jicong 已提交
1583 1584
  tstrncpy(pRspObj->topic, pWrapper->topicHandle->topicName, TSDB_TOPIC_FNAME_LEN);
  tstrncpy(pRspObj->db, pWrapper->topicHandle->db, TSDB_DB_FNAME_LEN);
1585

L
Liu Jicong 已提交
1586
  pRspObj->vgId = pWrapper->vgHandle->vgId;
L
Liu Jicong 已提交
1587
  pRspObj->resIter = -1;
L
Liu Jicong 已提交
1588
  memcpy(&pRspObj->rsp, &pWrapper->dataRsp, sizeof(SMqDataRsp));
L
Liu Jicong 已提交
1589

L
Liu Jicong 已提交
1590 1591
  pRspObj->resInfo.totalRows = 0;
  pRspObj->resInfo.precision = TSDB_TIME_PRECISION_MILLI;
1592

L
Liu Jicong 已提交
1593
  if (!pWrapper->dataRsp.withSchema) {
L
Liu Jicong 已提交
1594 1595
    setResSchemaInfo(&pRspObj->resInfo, pWrapper->topicHandle->schema.pSchema, pWrapper->topicHandle->schema.nCols);
  }
L
Liu Jicong 已提交
1596

1597 1598 1599 1600 1601
  // extract the rows in this data packet
  for(int32_t i = 0; i < pRspObj->rsp.blockNum; ++i) {
    SRetrieveTableRsp* pRetrieve = (SRetrieveTableRsp*)taosArrayGetP(pRspObj->rsp.blockData, i);
    int64_t rows = htobe64(pRetrieve->numOfRows);
    pVg->numOfRows += rows;
1602
    (*numOfRows) += rows;
1603 1604
  }

L
Liu Jicong 已提交
1605
  return pRspObj;
X
Xiaoyu Wang 已提交
1606 1607
}

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

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

  return pRspObj;
}

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 1655 1656 1657 1658
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 已提交
1659
  pParam->pVg = pVg;  // pVg may be released,fix it
1660 1661
  pParam->pTopic = pTopic;
  pParam->vgId = pVg->vgId;
H
Haojun Liao 已提交
1662
  pParam->requestId = req.reqId;
1663 1664 1665 1666 1667 1668 1669 1670 1671 1672 1673 1674 1675 1676 1677 1678 1679 1680 1681 1682 1683 1684 1685 1686

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

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

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

H
Haojun Liao 已提交
1687
  tscDebug("consumer:0x%" PRIx64 " send poll to %s vgId:%d, epoch %d, req:%s, reqId:0x%" PRIx64,
1688 1689 1690 1691 1692 1693 1694 1695 1696
           pTmq->consumerId, pTopic->topicName, pVg->vgId, pTmq->epoch, offsetFormatBuf, req.reqId);
  asyncSendMsgToServer(pTmq->pTscObj->pAppInfo->pTransporter, &pVg->epSet, &transporterId, sendInfo);

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

  return TSDB_CODE_SUCCESS;
}

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

  for (int i = 0; i < numOfTopics; i++) {
X
Xiaoyu Wang 已提交
1703
    SMqClientTopic* pTopic = taosArrayGet(tmq->clientTopics, i);
X
Xiaoyu Wang 已提交
1704
    int32_t         numOfVg = taosArrayGetSize(pTopic->vgs);
1705 1706

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

1714
      int32_t vgStatus = atomic_val_compare_exchange_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE, TMQ_VG_STATUS__WAIT);
1715
      if (vgStatus == TMQ_VG_STATUS__WAIT) {
L
Liu Jicong 已提交
1716
        int32_t vgSkipCnt = atomic_add_fetch_32(&pVg->vgSkipCnt, 1);
H
Haojun Liao 已提交
1717
        tscTrace("consumer:0x%" PRIx64 " epoch %d wait poll-rsp, skip vgId:%d skip cnt %d", tmq->consumerId, tmq->epoch,
X
Xiaoyu Wang 已提交
1718
                 pVg->vgId, vgSkipCnt);
X
Xiaoyu Wang 已提交
1719
        continue;
L
temp  
Liu Jicong 已提交
1720 1721 1722 1723
#if 0
        if (skipCnt < 30000) {
          continue;
        } else {
1724
        tscDebug("consumer:0x%" PRIx64 ",skip vgId:%d skip too much reset", tmq->consumerId, pVg->vgId);
L
temp  
Liu Jicong 已提交
1725 1726
        }
#endif
X
Xiaoyu Wang 已提交
1727
      }
1728

L
Liu Jicong 已提交
1729
      atomic_store_32(&pVg->vgSkipCnt, 0);
1730 1731 1732
      int32_t code = doTmqPollImpl(tmq, pTopic, pVg, timeout);
      if (code != TSDB_CODE_SUCCESS) {
        return code;
D
dapan1121 已提交
1733
      }
X
Xiaoyu Wang 已提交
1734 1735
    }
  }
1736

1737
  tscDebug("consumer:0x%" PRIx64 " end to poll data", tmq->consumerId);
X
Xiaoyu Wang 已提交
1738 1739 1740
  return 0;
}

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

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

X
Xiaoyu Wang 已提交
1764
  while (1) {
1765 1766
    SMqRspWrapper* pRspWrapper = NULL;
    taosGetQitem(tmq->qall, (void**)&pRspWrapper);
1767

1768
    if (pRspWrapper == NULL) {
L
Liu Jicong 已提交
1769
      taosReadAllQitems(tmq->mqueue, tmq->qall);
1770 1771
      taosGetQitem(tmq->qall, (void**)&pRspWrapper);
      if (pRspWrapper == NULL) {
L
Liu Jicong 已提交
1772 1773
        return NULL;
      }
X
Xiaoyu Wang 已提交
1774 1775
    }

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

1778 1779
    if (pRspWrapper->tmqRspType == TMQ_MSG_TYPE__END_RSP) {
      taosFreeQitem(pRspWrapper);
L
Liu Jicong 已提交
1780
      terrno = TSDB_CODE_TQ_NO_COMMITTED_OFFSET;
H
Haojun Liao 已提交
1781
      tscError("consumer:0x%" PRIx64 " unexpected rsp from poll, code:%s", tmq->consumerId, tstrerror(terrno));
L
Liu Jicong 已提交
1782
      return NULL;
1783 1784
    } else if (pRspWrapper->tmqRspType == TMQ_MSG_TYPE__POLL_RSP) {
      SMqPollRspWrapper* pollRspWrapper = (SMqPollRspWrapper*)pRspWrapper;
H
Haojun Liao 已提交
1785

1786
      int32_t consumerEpoch = atomic_load_32(&tmq->epoch);
1787 1788 1789
      SMqDataRsp* pDataRsp = &pollRspWrapper->dataRsp;

      if (pDataRsp->head.epoch == consumerEpoch) {
L
Liu Jicong 已提交
1790
        SMqClientVg* pVg = pollRspWrapper->vgHandle;
1791 1792 1793 1794 1795 1796 1797 1798 1799 1800

        // update the epset
        if (pollRspWrapper->pEpset != NULL) {
          SEp* pEp = GET_ACTIVE_EP(pollRspWrapper->pEpset);
          SEp* pOld = GET_ACTIVE_EP(&(pVg->epSet));
          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);
          pVg->epSet = *pollRspWrapper->pEpset;
        }

1801
        // update the local offset value only for the returned values.
1802
        pVg->currentOffset = pDataRsp->rspOffset;
X
Xiaoyu Wang 已提交
1803
        atomic_store_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE);
H
Haojun Liao 已提交
1804

1805 1806 1807
        char buf[80];
        tFormatOffset(buf, 80, &pDataRsp->rspOffset);
        if (pDataRsp->blockNum == 0) {
H
Haojun Liao 已提交
1808 1809
          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);
1810
          pRspWrapper = tmqFreeRspWrapper(pRspWrapper);
L
Liu Jicong 已提交
1811
          taosFreeQitem(pollRspWrapper);
1812
        } else {  // build rsp
1813 1814
          int64_t numOfRows = 0;
          SMqRspObj* pRsp = tmqBuildRspFromWrapper(pollRspWrapper, pVg, &numOfRows);
H
Haojun Liao 已提交
1815 1816
          tmq->totalRows += numOfRows;

H
Haojun Liao 已提交
1817
          tscDebug("consumer:0x%" PRIx64 " process poll rsp, vgId:%d, offset:%s, blocks:%d, rows:%" PRId64
H
Haojun Liao 已提交
1818
                   " vg total:%" PRId64 " total:%" PRId64 ", reqId:0x%" PRIx64,
H
Haojun Liao 已提交
1819
                   tmq->consumerId, pVg->vgId, buf, pDataRsp->blockNum, numOfRows, pVg->numOfRows, tmq->totalRows,
H
Haojun Liao 已提交
1820
                   pollRspWrapper->reqId);
1821 1822 1823
          taosFreeQitem(pollRspWrapper);
          return pRsp;
        }
X
Xiaoyu Wang 已提交
1824
      } else {
H
Haojun Liao 已提交
1825
        tscDebug("consumer:0x%" PRIx64 " vgId:%d msg discard since epoch mismatch: msg epoch %d, consumer epoch %d",
1826
                 tmq->consumerId, pollRspWrapper->vgId, pDataRsp->head.epoch, consumerEpoch);
1827
        pRspWrapper = tmqFreeRspWrapper(pRspWrapper);
L
Liu Jicong 已提交
1828 1829
        taosFreeQitem(pollRspWrapper);
      }
1830 1831
    } else if (pRspWrapper->tmqRspType == TMQ_MSG_TYPE__POLL_META_RSP) {
      SMqPollRspWrapper* pollRspWrapper = (SMqPollRspWrapper*)pRspWrapper;
L
Liu Jicong 已提交
1832
      int32_t            consumerEpoch = atomic_load_32(&tmq->epoch);
1833 1834 1835

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

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

L
Liu Jicong 已提交
1854 1855 1856 1857
      if (pollRspWrapper->taosxRsp.head.epoch == consumerEpoch) {
        SMqClientVg* pVg = pollRspWrapper->vgHandle;
        pVg->currentOffset = pollRspWrapper->taosxRsp.rspOffset;
        atomic_store_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE);
H
Haojun Liao 已提交
1858

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

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

1879 1880
        tmq->totalRows += numOfRows;

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

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

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

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

1911
TAOS_RES* tmq_consumer_poll(tmq_t* tmq, int64_t timeout) {
L
Liu Jicong 已提交
1912 1913
  void*   rspObj;
  int64_t startTime = taosGetTimestampMs();
L
Liu Jicong 已提交
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 1998 1999
int32_t tmq_consumer_close(tmq_t* tmq) {
  tscDebug("consumer:0x%" PRIx64" start to close consumer, status:%d", tmq->consumerId, tmq->status);
  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 2024
  } else {
    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 2129 2130 2131 2132 2133
typedef struct SSyncCommitInfo {
  tsem_t  sem;
  int32_t code;
} SSyncCommitInfo;

static void commitCallBackFn(tmq_t *pTmq, int32_t code, void* param) {
  SSyncCommitInfo* pInfo = (SSyncCommitInfo*) param;
  pInfo->code = code;
  tsem_post(&pInfo->sem);
2134
}
2135

2136 2137 2138 2139 2140 2141 2142 2143 2144
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 已提交
2145 2146
  } else {
    asyncCommitOffset(tmq, pRes, commitCallBackFn, pInfo);
2147 2148
  }

2149 2150
  tsem_wait(&pInfo->sem);
  code = pInfo->code;
H
Haojun Liao 已提交
2151 2152

  tsem_destroy(&pInfo->sem);
2153 2154 2155 2156 2157 2158 2159 2160 2161 2162 2163 2164 2165 2166 2167 2168 2169 2170 2171
  taosMemoryFree(pInfo);

  tscDebug("consumer:0x%"PRIx64" sync commit done, code:%s", tmq->consumerId, tstrerror(code));
  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);
    tmqUpdateEp(pTmq, head->epoch, &rsp);
    tDeleteSMqAskEpRsp(&rsp);
  }

H
Haojun Liao 已提交
2172
  tsem_post(&pInfo->sem);
2173 2174 2175 2176 2177 2178 2179 2180 2181 2182 2183 2184 2185 2186 2187 2188 2189 2190 2191 2192 2193 2194 2195 2196 2197 2198
}

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 已提交
2199
  tsem_init(&pInfo->sem, 0, 0);
2200 2201

  asyncAskEp(pTmq, updateEpCallbackFn, pInfo);
H
Haojun Liao 已提交
2202
  tsem_wait(&pInfo->sem);
2203 2204

  int32_t code = pInfo->code;
H
Haojun Liao 已提交
2205
  tsem_destroy(&pInfo->sem);
2206 2207 2208 2209 2210
  taosMemoryFree(pInfo);
  return code;
}

void asyncAskEp(tmq_t* pTmq, __tmq_askep_fn_t askEpFn, void* param) {
2211
  SMqAskEpReq req = {0};
2212 2213 2214
  req.consumerId = pTmq->consumerId;
  req.epoch = pTmq->epoch;
  strcpy(req.cgroup, pTmq->groupId);
2215 2216 2217

  int32_t tlen = tSerializeSMqAskEpReq(NULL, 0, &req);
  if (tlen < 0) {
2218 2219 2220
    tscError("consumer:0x%" PRIx64 ", tSerializeSMqAskEpReq failed", pTmq->consumerId);
    askEpFn(pTmq, TSDB_CODE_INVALID_PARA, NULL, param);
    return;
2221 2222 2223 2224
  }

  void* pReq = taosMemoryCalloc(1, tlen);
  if (pReq == NULL) {
2225 2226 2227
    tscError("consumer:0x%" PRIx64 ", failed to malloc askEpReq msg, size:%d", pTmq->consumerId, tlen);
    askEpFn(pTmq, TSDB_CODE_OUT_OF_MEMORY, NULL, param);
    return;
2228 2229 2230
  }

  if (tSerializeSMqAskEpReq(pReq, tlen, &req) < 0) {
2231
    tscError("consumer:0x%" PRIx64 ", tSerializeSMqAskEpReq %d failed", pTmq->consumerId, tlen);
2232
    taosMemoryFree(pReq);
2233 2234 2235

    askEpFn(pTmq, TSDB_CODE_INVALID_PARA, NULL, param);
    return;
2236 2237 2238 2239
  }

  SMqAskEpCbParam* pParam = taosMemoryCalloc(1, sizeof(SMqAskEpCbParam));
  if (pParam == NULL) {
2240
    tscError("consumer:0x%" PRIx64 ", failed to malloc subscribe param", pTmq->consumerId);
2241
    taosMemoryFree(pReq);
2242 2243 2244

    askEpFn(pTmq, TSDB_CODE_OUT_OF_MEMORY, NULL, param);
    return;
2245 2246
  }

2247 2248 2249 2250
  pParam->refId = pTmq->refId;
  pParam->epoch = pTmq->epoch;
  pParam->pUserFn = askEpFn;
  pParam->pParam = param;
2251 2252 2253 2254 2255

  SMsgSendInfo* sendInfo = taosMemoryCalloc(1, sizeof(SMsgSendInfo));
  if (sendInfo == NULL) {
    taosMemoryFree(pParam);
    taosMemoryFree(pReq);
2256 2257
    askEpFn(pTmq, TSDB_CODE_OUT_OF_MEMORY, NULL, param);
    return;
2258 2259
  }

2260
  sendInfo->msgInfo = (SDataBuf){ .pData = pReq, .len = tlen, .handle = NULL };
2261 2262 2263 2264

  sendInfo->requestId = generateRequestId();
  sendInfo->requestObjRefId = 0;
  sendInfo->param = pParam;
2265
  sendInfo->fp = askEpCallbackFn;
2266 2267
  sendInfo->msgType = TDMT_MND_TMQ_ASK_EP;

2268 2269
  SEpSet epSet = getEpSet_s(&pTmq->pTscObj->pAppInfo->mgmtEp);
  tscDebug("consumer:0x%" PRIx64 " ask ep from mnode, reqId:0x%" PRIx64, pTmq->consumerId, sendInfo->requestId);
2270 2271

  int64_t transporterId = 0;
2272
  asyncSendMsgToServer(pTmq->pTscObj->pAppInfo->pTransporter, &epSet, &transporterId, sendInfo);
2273 2274 2275 2276 2277 2278 2279
}

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

int32_t tmqCommitDone(SMqCommitCbParamSet* pParamSet) {
2280 2281 2282
  int64_t refId = pParamSet->refId;

  tmq_t* tmq = taosAcquireRef(tmqMgmt.rsetId, refId);
2283 2284 2285 2286 2287 2288 2289
  if (tmq == NULL) {
    taosMemoryFree(pParamSet);
    terrno = TSDB_CODE_TMQ_CONSUMER_CLOSED;
    return -1;
  }

  // if no more waiting rsp
2290
  pParamSet->callbackFn(tmq, pParamSet->code, pParamSet->userParam);
2291
  taosMemoryFree(pParamSet);
2292 2293

  taosReleaseRef(tmqMgmt.rsetId, refId);
2294
  return 0;
2295 2296
}

2297
void commitRspCountDown(SMqCommitCbParamSet* pParamSet, int64_t consumerId, const char* pTopic, int32_t vgId) {
2298 2299
  int32_t waitingRspNum = atomic_sub_fetch_32(&pParamSet->waitingRspNum, 1);
  if (waitingRspNum == 0) {
H
Haojun Liao 已提交
2300 2301
    tscDebug("consumer:0x%" PRIx64 " topic:%s vgId:%d all commit-rsp received, commit completed", consumerId, pTopic,
             vgId);
2302
    tmqCommitDone(pParamSet);
H
Haojun Liao 已提交
2303 2304 2305
  } else {
    tscDebug("consumer:0x%" PRIx64 " topic:%s vgId:%d commit-rsp received, remain:%d", consumerId, pTopic, vgId,
             waitingRspNum);
2306 2307
  }
}