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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

typedef struct {
  int64_t       pollCnt;
  int64_t       numOfRows;
  SVgOffsetInfo offsetInfo;
  int32_t       vgId;
  int32_t       vgStatus;
  int32_t       vgSkipCnt;            // here used to mark the slow vgroups
  int64_t       emptyBlockReceiveTs;  // once empty block is received, idle for ignoreCnt then start to poll data
  SEpSet        epSet;
152 153
} SMqClientVg;

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

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

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

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

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

199
typedef struct {
200 201
  int64_t        refId;
  int32_t        epoch;
L
Liu Jicong 已提交
202 203
  int32_t        waitingRspNum;
  int32_t        totalRspNum;
204
  int32_t        code;
205
  tmq_commit_cb* callbackFn;
L
Liu Jicong 已提交
206 207
  /*SArray*        successfulOffsets;*/
  /*SArray*        failedOffsets;*/
X
Xiaoyu Wang 已提交
208
  void* userParam;
209 210 211 212 213
} SMqCommitCbParamSet;

typedef struct {
  SMqCommitCbParamSet* params;
  STqOffset*           pOffset;
H
Haojun Liao 已提交
214 215 216
  char                 topicName[TSDB_TOPIC_FNAME_LEN];
  int32_t              vgId;
  tmq_t*               pTmq;
217
} SMqCommitCbParam;
218

219 220 221 222 223
typedef struct SSyncCommitInfo {
  tsem_t  sem;
  int32_t code;
} SSyncCommitInfo;

224
static int32_t doAskEp(tmq_t* tmq);
225 226
static int32_t makeTopicVgroupKey(char* dst, const char* topicName, int32_t vg);
static int32_t tmqCommitDone(SMqCommitCbParamSet* pParamSet);
227 228
static int32_t doSendCommitMsg(tmq_t* tmq, SMqClientVg* pVg, const char* pTopicName, SMqCommitCbParamSet* pParamSet,
                               int32_t index, int32_t totalVgroups);
X
Xiaoyu Wang 已提交
229 230 231
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);
232

233
tmq_conf_t* tmq_conf_new() {
wafwerar's avatar
wafwerar 已提交
234
  tmq_conf_t* conf = taosMemoryCalloc(1, sizeof(tmq_conf_t));
235 236 237 238 239
  if (conf == NULL) {
    terrno = TSDB_CODE_OUT_OF_MEMORY;
    return conf;
  }

240
  conf->withTbName = false;
L
Liu Jicong 已提交
241
  conf->autoCommit = true;
242
  conf->autoCommitInterval = DEFAULT_AUTO_COMMIT_INTERVAL;
243
  conf->resetOffset = TMQ_OFFSET__RESET_EARLIEAST;
244
  conf->hbBgEnable = true;
245

246 247 248
  return conf;
}

L
Liu Jicong 已提交
249
void tmq_conf_destroy(tmq_conf_t* conf) {
L
Liu Jicong 已提交
250
  if (conf) {
251 252 253 254 255 256 257 258 259
    if (conf->ip) {
      taosMemoryFree(conf->ip);
    }
    if (conf->user) {
      taosMemoryFree(conf->user);
    }
    if (conf->pass) {
      taosMemoryFree(conf->pass);
    }
L
Liu Jicong 已提交
260 261
    taosMemoryFree(conf);
  }
L
Liu Jicong 已提交
262 263 264
}

tmq_conf_res_t tmq_conf_set(tmq_conf_t* conf, const char* key, const char* value) {
265
  if (strcasecmp(key, "group.id") == 0) {
L
Liu Jicong 已提交
266
    tstrncpy(conf->groupId, value, TSDB_CGROUP_LEN);
L
Liu Jicong 已提交
267
    return TMQ_CONF_OK;
268
  }
L
Liu Jicong 已提交
269

270
  if (strcasecmp(key, "client.id") == 0) {
L
Liu Jicong 已提交
271
    tstrncpy(conf->clientId, value, 256);
L
Liu Jicong 已提交
272 273
    return TMQ_CONF_OK;
  }
L
Liu Jicong 已提交
274

275 276
  if (strcasecmp(key, "enable.auto.commit") == 0) {
    if (strcasecmp(value, "true") == 0) {
L
Liu Jicong 已提交
277
      conf->autoCommit = true;
L
Liu Jicong 已提交
278
      return TMQ_CONF_OK;
279
    } else if (strcasecmp(value, "false") == 0) {
L
Liu Jicong 已提交
280
      conf->autoCommit = false;
L
Liu Jicong 已提交
281 282 283 284
      return TMQ_CONF_OK;
    } else {
      return TMQ_CONF_INVALID;
    }
285
  }
L
Liu Jicong 已提交
286

287
  if (strcasecmp(key, "auto.commit.interval.ms") == 0) {
288
    conf->autoCommitInterval = taosStr2int64(value);
L
Liu Jicong 已提交
289 290 291
    return TMQ_CONF_OK;
  }

292 293 294
  if (strcasecmp(key, "auto.offset.reset") == 0) {
    if (strcasecmp(value, "none") == 0) {
      conf->resetOffset = TMQ_OFFSET__RESET_NONE;
L
Liu Jicong 已提交
295
      return TMQ_CONF_OK;
296 297
    } else if (strcasecmp(value, "earliest") == 0) {
      conf->resetOffset = TMQ_OFFSET__RESET_EARLIEAST;
L
Liu Jicong 已提交
298
      return TMQ_CONF_OK;
299 300
    } else if (strcasecmp(value, "latest") == 0) {
      conf->resetOffset = TMQ_OFFSET__RESET_LATEST;
L
Liu Jicong 已提交
301 302 303 304 305
      return TMQ_CONF_OK;
    } else {
      return TMQ_CONF_INVALID;
    }
  }
L
Liu Jicong 已提交
306

307 308
  if (strcasecmp(key, "msg.with.table.name") == 0) {
    if (strcasecmp(value, "true") == 0) {
309
      conf->withTbName = true;
L
Liu Jicong 已提交
310
      return TMQ_CONF_OK;
311
    } else if (strcasecmp(value, "false") == 0) {
312
      conf->withTbName = false;
L
Liu Jicong 已提交
313
      return TMQ_CONF_OK;
314 315 316 317 318
    } else {
      return TMQ_CONF_INVALID;
    }
  }

319 320
  if (strcasecmp(key, "experimental.snapshot.enable") == 0) {
    if (strcasecmp(value, "true") == 0) {
L
Liu Jicong 已提交
321
      conf->snapEnable = true;
322
      return TMQ_CONF_OK;
323
    } else if (strcasecmp(value, "false") == 0) {
L
Liu Jicong 已提交
324
      conf->snapEnable = false;
325 326 327 328 329 330
      return TMQ_CONF_OK;
    } else {
      return TMQ_CONF_INVALID;
    }
  }

331
  if (strcasecmp(key, "experimental.snapshot.batch.size") == 0) {
332
    conf->snapBatchSize = taosStr2int64(value);
L
Liu Jicong 已提交
333 334 335
    return TMQ_CONF_OK;
  }

336
  if (strcasecmp(key, "enable.heartbeat.background") == 0) {
X
Xiaoyu Wang 已提交
337 338 339 340 341 342 343 344 345 346
    //    if (strcasecmp(value, "true") == 0) {
    //      conf->hbBgEnable = true;
    //      return TMQ_CONF_OK;
    //    } else if (strcasecmp(value, "false") == 0) {
    //      conf->hbBgEnable = false;
    //      return TMQ_CONF_OK;
    //    } else {
    tscError("the default value of enable.heartbeat.background is true, can not be seted");
    return TMQ_CONF_INVALID;
    //    }
L
Liu Jicong 已提交
347 348
  }

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

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

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

364
  if (strcasecmp(key, "td.connect.port") == 0) {
365
    conf->port = taosStr2int64(value);
L
Liu Jicong 已提交
366 367
    return TMQ_CONF_OK;
  }
368

369
  if (strcasecmp(key, "td.connect.db") == 0) {
L
Liu Jicong 已提交
370 371 372
    return TMQ_CONF_OK;
  }

L
Liu Jicong 已提交
373
  return TMQ_CONF_UNKNOWN;
374 375
}

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

L
Liu Jicong 已提交
378 379
int32_t tmq_list_append(tmq_list_t* list, const char* src) {
  SArray* container = &list->container;
380
  if (src == NULL || src[0] == 0) return -1;
381
  char* topic = taosStrdup(src);
382 383 384
  if (topic[0] != '`') {
    strtolower(topic, src);
  }
L
fix  
Liu Jicong 已提交
385
  if (taosArrayPush(container, &topic) == NULL) return -1;
386 387 388
  return 0;
}

L
Liu Jicong 已提交
389
void tmq_list_destroy(tmq_list_t* list) {
L
Liu Jicong 已提交
390
  SArray* container = &list->container;
L
Liu Jicong 已提交
391
  taosArrayDestroyP(container, taosMemoryFree);
L
Liu Jicong 已提交
392 393
}

L
Liu Jicong 已提交
394 395 396 397 398 399 400 401 402 403
int32_t tmq_list_get_size(const tmq_list_t* list) {
  const SArray* container = &list->container;
  return taosArrayGetSize(container);
}

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

X
Xiaoyu Wang 已提交
404 405
static SMqClientVg* foundClientVg(SArray* pTopicList, const char* pName, int32_t vgId, int32_t* index,
                                  int32_t* numOfVgroups) {
406 407 408 409
  int32_t numOfTopics = taosArrayGetSize(pTopicList);
  *index = -1;
  *numOfVgroups = 0;

X
Xiaoyu Wang 已提交
410
  for (int32_t i = 0; i < numOfTopics; ++i) {
411 412
    SMqClientTopic* pTopic = taosArrayGet(pTopicList, i);
    if (strcmp(pTopic->topicName, pName) != 0) {
413 414 415
      continue;
    }

416 417
    *numOfVgroups = taosArrayGetSize(pTopic->vgs);
    for (int32_t j = 0; j < (*numOfVgroups); ++j) {
418
      SMqClientVg* pClientVg = taosArrayGet(pTopic->vgs, j);
419 420 421
      if (pClientVg->vgId == vgId) {
        *index = j;
        return pClientVg;
422 423
      }
    }
L
Liu Jicong 已提交
424
  }
425 426

  return NULL;
L
Liu Jicong 已提交
427
}
428

429 430 431
// 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 已提交
432
static int32_t tmqCommitCb(void* param, SDataBuf* pBuf, int32_t code) {
433
  SMqCommitCbParam*    pParam = (SMqCommitCbParam*)param;
434
  SMqCommitCbParamSet* pParamSet = (SMqCommitCbParamSet*)pParam->params;
H
Haojun Liao 已提交
435

X
Xiaoyu Wang 已提交
436 437 438 439 440 441 442 443 444 445 446 447 448 449
  //  if (code != TSDB_CODE_SUCCESS) { // if commit offset failed, let's try again
  //    taosThreadMutexLock(&pParam->pTmq->lock);
  //    int32_t numOfVgroups, index;
  //    SMqClientVg* pVg = foundClientVg(pParam->pTmq->clientTopics, pParam->topicName, pParam->vgId, &index,
  //    &numOfVgroups); if (pVg == NULL) {
  //      tscDebug("consumer:0x%" PRIx64
  //               " subKey:%s vgId:%d commit failed, code:%s has been transferred to other consumer, no need retry
  //               ordinal:%d/%d", pParam->pTmq->consumerId, pParam->pOffset->subKey, pParam->vgId, tstrerror(code),
  //               index + 1, numOfVgroups);
  //    } else { // let's retry the commit
  //      int32_t code1 = doSendCommitMsg(pParam->pTmq, pVg, pParam->topicName, pParamSet, index, numOfVgroups);
  //      if (code1 != TSDB_CODE_SUCCESS) {  // retry failed.
  //        tscError("consumer:0x%" PRIx64 " topic:%s vgId:%d offset:%" PRId64
  //                 " retry failed, ignore this commit. code:%s ordinal:%d/%d",
H
Haojun Liao 已提交
450
  //                 pParam->pTmq->consumerId, pParam->topicName, pVg->vgId, pVg->offsetInfo.committedOffset.version,
X
Xiaoyu Wang 已提交
451 452 453 454 455 456 457 458 459 460 461 462 463 464 465
  //                 tstrerror(terrno), index + 1, numOfVgroups);
  //      }
  //    }
  //
  //    taosThreadMutexUnlock(&pParam->pTmq->lock);
  //
  //    taosMemoryFree(pParam->pOffset);
  //    taosMemoryFree(pBuf->pData);
  //    taosMemoryFree(pBuf->pEpSet);
  //
  //    commitRspCountDown(pParamSet, pParam->pTmq->consumerId, pParam->topicName, pParam->vgId);
  //    return 0;
  //  }
  //
  //  // todo replace the pTmq with refId
466

L
Liu Jicong 已提交
467
  taosMemoryFree(pParam->pOffset);
L
Liu Jicong 已提交
468
  taosMemoryFree(pBuf->pData);
dengyihao's avatar
dengyihao 已提交
469
  taosMemoryFree(pBuf->pEpSet);
L
Liu Jicong 已提交
470

471
  commitRspCountDown(pParamSet, pParam->pTmq->consumerId, pParam->topicName, pParam->vgId);
472 473 474
  return 0;
}

475 476
static int32_t doSendCommitMsg(tmq_t* tmq, SMqClientVg* pVg, const char* pTopicName, SMqCommitCbParamSet* pParamSet,
                               int32_t index, int32_t totalVgroups) {
L
Liu Jicong 已提交
477 478
  STqOffset* pOffset = taosMemoryCalloc(1, sizeof(STqOffset));
  if (pOffset == NULL) {
479
    return TSDB_CODE_OUT_OF_MEMORY;
L
Liu Jicong 已提交
480
  }
481

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

L
Liu Jicong 已提交
484 485 486
  int32_t groupLen = strlen(tmq->groupId);
  memcpy(pOffset->subKey, tmq->groupId, groupLen);
  pOffset->subKey[groupLen] = TMQ_SEPARATOR;
H
Haojun Liao 已提交
487
  strcpy(pOffset->subKey + groupLen + 1, pTopicName);
L
Liu Jicong 已提交
488

489 490
  int32_t len = 0;
  int32_t code = 0;
L
Liu Jicong 已提交
491 492
  tEncodeSize(tEncodeSTqOffset, pOffset, len, code);
  if (code < 0) {
493
    return TSDB_CODE_INVALID_PARA;
L
Liu Jicong 已提交
494
  }
495

L
Liu Jicong 已提交
496
  void* buf = taosMemoryCalloc(1, sizeof(SMsgHead) + len);
L
Liu Jicong 已提交
497 498
  if (buf == NULL) {
    taosMemoryFree(pOffset);
499
    return TSDB_CODE_OUT_OF_MEMORY;
L
Liu Jicong 已提交
500
  }
501

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

L
Liu Jicong 已提交
504 505 506 507 508
  void* abuf = POINTER_SHIFT(buf, sizeof(SMsgHead));

  SEncoder encoder;
  tEncoderInit(&encoder, abuf, len);
  tEncodeSTqOffset(&encoder, pOffset);
L
Liu Jicong 已提交
509
  tEncoderClear(&encoder);
L
Liu Jicong 已提交
510 511

  // build param
512
  SMqCommitCbParam* pParam = taosMemoryCalloc(1, sizeof(SMqCommitCbParam));
L
Liu Jicong 已提交
513
  if (pParam == NULL) {
L
Liu Jicong 已提交
514
    taosMemoryFree(pOffset);
L
Liu Jicong 已提交
515
    taosMemoryFree(buf);
516
    return TSDB_CODE_OUT_OF_MEMORY;
L
Liu Jicong 已提交
517
  }
518

L
Liu Jicong 已提交
519 520
  pParam->params = pParamSet;
  pParam->pOffset = pOffset;
H
Haojun Liao 已提交
521 522 523
  pParam->vgId = pVg->vgId;
  pParam->pTmq = tmq;

H
Haojun Liao 已提交
524
  tstrncpy(pParam->topicName, pTopicName, tListLen(pParam->topicName));
L
Liu Jicong 已提交
525 526 527 528

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

535
  pMsgSendInfo->msgInfo = (SDataBuf) { .pData = buf, .len = sizeof(SMsgHead) + len, .handle = NULL };
L
Liu Jicong 已提交
536 537 538 539

  pMsgSendInfo->requestId = generateRequestId();
  pMsgSendInfo->requestObjRefId = 0;
  pMsgSendInfo->param = pParam;
L
Liu Jicong 已提交
540
  pMsgSendInfo->paramFreeFp = taosMemoryFree;
541
  pMsgSendInfo->fp = tmqCommitCb;
L
Liu Jicong 已提交
542
  pMsgSendInfo->msgType = TDMT_VND_TMQ_COMMIT_OFFSET;
L
Liu Jicong 已提交
543

L
Liu Jicong 已提交
544 545 546
  atomic_add_fetch_32(&pParamSet->waitingRspNum, 1);
  atomic_add_fetch_32(&pParamSet->totalRspNum, 1);

H
Haojun Liao 已提交
547
  SEp* pEp = GET_ACTIVE_EP(&pVg->epSet);
548 549 550 551
  char offsetBuf[80] = {0};
  tFormatOffset(offsetBuf, tListLen(offsetBuf), &pOffset->val);

  char commitBuf[80] = {0};
H
Haojun Liao 已提交
552
  tFormatOffset(commitBuf, tListLen(commitBuf), &pVg->offsetInfo.committedOffset);
553 554 555
  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 已提交
556

L
Liu Jicong 已提交
557 558
  int64_t transporterId = 0;
  asyncSendMsgToServer(tmq->pTscObj->pAppInfo->pTransporter, &pVg->epSet, &transporterId, pMsgSendInfo);
559 560

  return TSDB_CODE_SUCCESS;
L
Liu Jicong 已提交
561 562
}

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

  SMqCommitCbParamSet* pParamSet = taosMemoryCalloc(1, sizeof(SMqCommitCbParamSet));
  if (pParamSet == NULL) {
592 593
    pCommitFp(tmq, TSDB_CODE_OUT_OF_MEMORY, userParam);
    return;
L
Liu Jicong 已提交
594
  }
H
Haojun Liao 已提交
595

596 597
  pParamSet->refId = tmq->refId;
  pParamSet->epoch = tmq->epoch;
598
  pParamSet->callbackFn = pCommitFp;
L
Liu Jicong 已提交
599
  pParamSet->userParam = userParam;
L
Liu Jicong 已提交
600

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

603 604 605 606
  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 已提交
607
    SMqClientTopic* pTopic = taosArrayGet(tmq->clientTopics, i);
608 609
    if (strcmp(pTopic->topicName, pTopicName) == 0) {
      break;
610
    }
611
  }
612

613
  if (i == numOfTopics) {
X
Xiaoyu Wang 已提交
614 615
    tscWarn("consumer:0x%" PRIx64 " failed to find the specified topic:%s, total topics:%d", tmq->consumerId,
            pTopicName, numOfTopics);
616 617 618 619
    taosMemoryFree(pParamSet);
    pCommitFp(tmq, TSDB_CODE_SUCCESS, userParam);
    return;
  }
L
Liu Jicong 已提交
620

621 622 623 624 625 626 627 628
  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 已提交
629
    }
L
Liu Jicong 已提交
630
  }
L
Liu Jicong 已提交
631

632
  if (j == numOfVgroups) {
X
Xiaoyu Wang 已提交
633 634
    tscWarn("consumer:0x%" PRIx64 " failed to find the specified vgId:%d, total Vgs:%d, topic:%s", tmq->consumerId,
            vgId, numOfVgroups, pTopicName);
L
Liu Jicong 已提交
635
    taosMemoryFree(pParamSet);
636 637
    pCommitFp(tmq, TSDB_CODE_SUCCESS, userParam);
    return;
L
Liu Jicong 已提交
638 639
  }

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

644 645 646 647 648
    // failed to commit, callback user function directly.
    if (code != TSDB_CODE_SUCCESS) {
      taosMemoryFree(pParamSet);
      pCommitFp(tmq, code, userParam);
    }
X
Xiaoyu Wang 已提交
649
  } else {  // do not perform commit, callback user function directly.
650 651
    taosMemoryFree(pParamSet);
    pCommitFp(tmq, code, userParam);
L
Liu Jicong 已提交
652 653 654
  }
}

655
static void asyncCommitAllOffsets(tmq_t* tmq, tmq_commit_cb* pCommitFp, void* userParam) {
656 657
  SMqCommitCbParamSet* pParamSet = taosMemoryCalloc(1, sizeof(SMqCommitCbParamSet));
  if (pParamSet == NULL) {
658 659
    pCommitFp(tmq, TSDB_CODE_OUT_OF_MEMORY, userParam);
    return;
660
  }
661 662 663

  pParamSet->refId = tmq->refId;
  pParamSet->epoch = tmq->epoch;
664
  pParamSet->callbackFn = pCommitFp;
665 666
  pParamSet->userParam = userParam;

667 668 669
  // init as 1 to prevent concurrency issue
  pParamSet->waitingRspNum = 1;

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

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

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

H
Haojun Liao 已提交
682
      if (pVg->offsetInfo.currentOffset.type > 0 && !tOffsetEqual(&pVg->offsetInfo.currentOffset, &pVg->offsetInfo.committedOffset)) {
683
        int32_t code = doSendCommitMsg(tmq, pVg, pTopic->topicName, pParamSet, j, numOfVgroups);
684 685
        if (code != TSDB_CODE_SUCCESS) {
          tscError("consumer:0x%" PRIx64 " topic:%s vgId:%d offset:%" PRId64 " failed, code:%s ordinal:%d/%d",
H
Haojun Liao 已提交
686
                   tmq->consumerId, pTopic->topicName, pVg->vgId, pVg->offsetInfo.committedOffset.version, tstrerror(terrno),
687
                   j + 1, numOfVgroups);
L
Liu Jicong 已提交
688 689
          continue;
        }
H
Haojun Liao 已提交
690 691

        // update the offset value.
H
Haojun Liao 已提交
692
        pVg->offsetInfo.committedOffset = pVg->offsetInfo.currentOffset;
693
      } else {
D
dapan1121 已提交
694
        tscDebug("consumer:0x%" PRIx64 " topic:%s vgId:%d, no commit, current:%" PRId64 ", ordinal:%d/%d",
H
Haojun Liao 已提交
695
                 tmq->consumerId, pTopic->topicName, pVg->vgId, pVg->offsetInfo.currentOffset.version, j + 1, numOfVgroups);
696 697 698 699
      }
    }
  }

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

L
Liu Jicong 已提交
703
  // no request is sent
L
Liu Jicong 已提交
704 705
  if (pParamSet->totalRspNum == 0) {
    taosMemoryFree(pParamSet);
706 707
    pCommitFp(tmq, TSDB_CODE_SUCCESS, userParam);
    return;
L
Liu Jicong 已提交
708 709
  }

L
Liu Jicong 已提交
710
  // count down since waiting rsp num init as 1
711
  commitRspCountDown(pParamSet, tmq->consumerId, "", 0);
712 713
}

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

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

void tmqAssignDelayedCommitTask(void* param, void* tmrId) {
732
  int64_t refId = *(int64_t*)param;
733
  generateTimedTask(refId, TMQ_DELAYED_TASK__COMMIT);
734
  taosMemoryFree(param);
L
Liu Jicong 已提交
735 736 737
}

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

  taosReleaseRef(tmqMgmt.rsetId, refId);
748
  taosMemoryFree(param);
L
Liu Jicong 已提交
749 750
}

751
int32_t tmqHbCb(void* param, SDataBuf* pMsg, int32_t code) {
dengyihao's avatar
dengyihao 已提交
752 753 754 755
  if (pMsg) {
    taosMemoryFree(pMsg->pData);
    taosMemoryFree(pMsg->pEpSet);
  }
756 757 758 759
  return 0;
}

void tmqSendHbReq(void* param, void* tmrId) {
760
  int64_t refId = *(int64_t*)param;
761

X
Xiaoyu Wang 已提交
762
  tmq_t* tmq = taosAcquireRef(tmqMgmt.rsetId, refId);
763
  if (tmq == NULL) {
L
Liu Jicong 已提交
764
    taosMemoryFree(param);
765 766
    return;
  }
D
dapan1121 已提交
767 768 769 770 771

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

L
Liu Jicong 已提交
772
  int32_t tlen = tSerializeSMqHbReq(NULL, 0, &req);
D
dapan1121 已提交
773 774
  if (tlen < 0) {
    tscError("tSerializeSMqHbReq failed");
775
    goto OVER;
D
dapan1121 已提交
776
  }
777

L
Liu Jicong 已提交
778
  void* pReq = taosMemoryCalloc(1, tlen);
D
dapan1121 已提交
779 780
  if (tlen < 0) {
    tscError("failed to malloc MqHbReq msg, size:%d", tlen);
781
    goto OVER;
D
dapan1121 已提交
782
  }
783

D
dapan1121 已提交
784 785 786
  if (tSerializeSMqHbReq(pReq, tlen, &req) < 0) {
    tscError("tSerializeSMqHbReq %d failed", tlen);
    taosMemoryFree(pReq);
787
    goto OVER;
D
dapan1121 已提交
788
  }
789 790 791 792

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

796
  sendInfo->msgInfo = (SDataBuf){ .pData = pReq, .len = tlen, .handle = NULL };
797 798 799 800 801

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

  taosFreeQall(tmq->qall);
982
  tsem_destroy(&tmq->rspSem);
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

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

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

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

1050 1051
  pTmq->hbBgEnable = conf->hbBgEnable;

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

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

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

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

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

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

1089
  return pTmq;
1090

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

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

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

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

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

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

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

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

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

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

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

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

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

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

1156
  if (tsem_init(&param.rspSem, 0, 0) != 0) {
wmmhello's avatar
wmmhello 已提交
1157
    code = TSDB_CODE_TSC_INTERNAL_ERROR;
1158 1159
    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) {
wmmhello's avatar
wmmhello 已提交
1193 1194
      tscError("consumer:0x%" PRIx64 ", mnd not ready for subscribe, retry:%d in 500ms", tmq->consumerId, retryCnt);
      code = TSDB_CODE_TSC_INTERNAL_ERROR;
L
Liu Jicong 已提交
1195 1196
      goto FAIL;
    }
1197

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

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

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

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

L
Liu Jicong 已提交
1221
  return code;
1222 1223
}

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

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

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

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

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

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

L
Liu Jicong 已提交
1251
  if (code != 0) {
H
Haojun Liao 已提交
1252 1253
    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 已提交
1254

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

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

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

L
fix txn  
Liu Jicong 已提交
1278
    goto CREATE_MSG_FAIL;
L
Liu Jicong 已提交
1279 1280
  }

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

1289
    tsem_post(&tmq->rspSem);
1290 1291
    taosReleaseRef(tmqMgmt.rsetId, refId);

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

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

L
Liu Jicong 已提交
1302 1303 1304
  // handle meta rsp
  int8_t rspType = ((SMqRspHead*)pMsg->pData)->mqMsgType;

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

L
Liu Jicong 已提交
1314
  pRspWrapper->tmqRspType = rspType;
L
Liu Jicong 已提交
1315 1316
  pRspWrapper->vgHandle = pVg;
  pRspWrapper->topicHandle = pTopic;
H
Haojun Liao 已提交
1317
  pRspWrapper->reqId = requestId;
1318
  pRspWrapper->pEpset = pMsg->pEpSet;
1319
  pRspWrapper->vgId = pVg->vgId;
L
Liu Jicong 已提交
1320

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

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

L
Liu Jicong 已提交
1349
  taosMemoryFree(pMsg->pData);
H
Haojun Liao 已提交
1350
  taosWriteQitem(tmq->mqueue, pRspWrapper);
L
Liu Jicong 已提交
1351

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

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

L
Liu Jicong 已提交
1359
  return 0;
H
Haojun Liao 已提交
1360

L
fix txn  
Liu Jicong 已提交
1361
CREATE_MSG_FAIL:
L
Liu Jicong 已提交
1362
  if (epoch == tmq->epoch) {
L
Liu Jicong 已提交
1363 1364
    atomic_store_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE);
  }
H
Haojun Liao 已提交
1365

1366
  tsem_post(&tmq->rspSem);
1367 1368
  taosReleaseRef(tmqMgmt.rsetId, refId);

L
Liu Jicong 已提交
1369
  return -1;
1370 1371
}

H
Haojun Liao 已提交
1372 1373 1374 1375 1376
typedef struct SVgroupSaveInfo {
  STqOffsetVal offset;
  int64_t      numOfRows;
} SVgroupSaveInfo;

H
Haojun Liao 已提交
1377 1378 1379 1380 1381 1382
static void initClientTopicFromRsp(SMqClientTopic* pTopic, SMqSubTopicEp* pTopicEp, SHashObj* pVgOffsetHashMap,
                                   tmq_t* tmq) {
  pTopic->schema = pTopicEp->schema;
  pTopicEp->schema.nCols = 0;
  pTopicEp->schema.pSchema = NULL;

X
Xiaoyu Wang 已提交
1383
  char    vgKey[TSDB_TOPIC_FNAME_LEN + 22];
H
Haojun Liao 已提交
1384 1385 1386 1387 1388 1389 1390 1391 1392 1393
  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 已提交
1394 1395

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

X
Xiaoyu Wang 已提交
1398 1399
    int64_t      numOfRows = 0;
    STqOffsetVal offsetNew = {.type = tmq->resetOffsetCfg};
H
Haojun Liao 已提交
1400 1401 1402
    if (pInfo != NULL) {
      offsetNew = pInfo->offset;
      numOfRows = pInfo->numOfRows;
H
Haojun Liao 已提交
1403 1404 1405 1406 1407 1408 1409 1410
    }

    SMqClientVg clientVg = {
        .pollCnt = 0,
        .vgId = pVgEp->vgId,
        .epSet = pVgEp->epSet,
        .vgStatus = TMQ_VG_STATUS__IDLE,
        .vgSkipCnt = 0,
H
Haojun Liao 已提交
1411
        .emptyBlockReceiveTs = 0,
H
Haojun Liao 已提交
1412
        .numOfRows = numOfRows,
H
Haojun Liao 已提交
1413 1414
    };

H
Haojun Liao 已提交
1415 1416 1417 1418
    clientVg.offsetInfo.currentOffset = offsetNew;
    clientVg.offsetInfo.committedOffset = offsetNew;
    clientVg.offsetInfo.walVerBegin = -1;
    clientVg.offsetInfo.walVerEnd = -1;
H
Haojun Liao 已提交
1419 1420 1421 1422 1423 1424 1425 1426 1427 1428 1429 1430 1431
    taosArrayPush(pTopic->vgs, &clientVg);
  }
}

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

  taosArrayDestroy(pTopic->vgs);
}

1432
static bool doUpdateLocalEp(tmq_t* tmq, int32_t epoch, const SMqAskEpRsp* pRsp) {
1433 1434
  bool set = false;

1435
  int32_t topicNumCur = taosArrayGetSize(tmq->clientTopics);
1436
  int32_t topicNumGet = taosArrayGetSize(pRsp->topics);
1437

X
Xiaoyu Wang 已提交
1438 1439
  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",
1440
           tmq->consumerId, tmq->epoch, epoch, topicNumGet, topicNumCur);
wmmhello's avatar
wmmhello 已提交
1441 1442 1443
  if (epoch <= tmq->epoch) {
    return false;
  }
1444 1445 1446 1447 1448 1449

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

H
Haojun Liao 已提交
1450 1451
  SHashObj* pVgOffsetHashMap = taosHashInit(64, MurmurHash3_32, false, HASH_NO_LOCK);
  if (pVgOffsetHashMap == NULL) {
1452 1453 1454
    taosArrayDestroy(newTopics);
    return false;
  }
1455

H
Haojun Liao 已提交
1456
  // todo extract method
1457 1458 1459 1460 1461
  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);
1462
      tscDebug("consumer:0x%" PRIx64 ", new vg num: %d", tmq->consumerId, vgNumCur);
1463 1464
      for (int32_t j = 0; j < vgNumCur; j++) {
        SMqClientVg* pVgCur = taosArrayGet(pTopicCur->vgs, j);
H
Haojun Liao 已提交
1465 1466
        makeTopicVgroupKey(vgKey, pTopicCur->topicName, pVgCur->vgId);

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

H
Haojun Liao 已提交
1472
        SVgroupSaveInfo info = {.offset = pVgCur->offsetInfo.currentOffset, .numOfRows = pVgCur->numOfRows};
H
Haojun Liao 已提交
1473
        taosHashPut(pVgOffsetHashMap, vgKey, strlen(vgKey), &info, sizeof(SVgroupSaveInfo));
1474 1475 1476 1477 1478 1479 1480
      }
    }
  }

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

H
Haojun Liao 已提交
1485 1486
  taosHashCleanup(pVgOffsetHashMap);

1487
  // destroy current buffered existed topics info
1488
  if (tmq->clientTopics) {
H
Haojun Liao 已提交
1489
    taosArrayDestroyEx(tmq->clientTopics, freeClientVgInfo);
X
Xiaoyu Wang 已提交
1490
  }
H
Haojun Liao 已提交
1491
  tmq->clientTopics = newTopics;
1492

X
Xiaoyu Wang 已提交
1493
  int8_t flag = (topicNumGet == 0) ? TMQ_CONSUMER_STATUS__NO_TOPIC : TMQ_CONSUMER_STATUS__READY;
H
Haojun Liao 已提交
1494
  atomic_store_8(&tmq->status, flag);
X
Xiaoyu Wang 已提交
1495
  atomic_store_32(&tmq->epoch, epoch);
H
Haojun Liao 已提交
1496

1497
  tscDebug("consumer:0x%" PRIx64 " update topic info completed", tmq->consumerId);
X
Xiaoyu Wang 已提交
1498 1499 1500
  return set;
}

1501
int32_t askEpCallbackFn(void* param, SDataBuf* pMsg, int32_t code) {
1502
  SMqAskEpCbParam* pParam = (SMqAskEpCbParam*)param;
1503 1504 1505
  tmq_t*           tmq = taosAcquireRef(tmqMgmt.rsetId, pParam->refId);

  if (tmq == NULL) {
1506 1507 1508
    terrno = TSDB_CODE_TMQ_CONSUMER_CLOSED;
    pParam->pUserFn(tmq, terrno, NULL, pParam->pParam);

1509
    taosMemoryFree(pMsg->pData);
dengyihao's avatar
dengyihao 已提交
1510
    taosMemoryFree(pMsg->pEpSet);
1511 1512
    taosMemoryFree(pParam);
    return terrno;
1513 1514
  }

H
Haojun Liao 已提交
1515
  if (code != TSDB_CODE_SUCCESS) {
1516 1517 1518 1519 1520 1521 1522 1523 1524
    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;
1525
  }
L
Liu Jicong 已提交
1526

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

1536 1537 1538 1539 1540 1541 1542 1543
    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 已提交
1544
  } else {
1545 1546
    tscDebug("consumer:0x%" PRIx64 ", recv ep, msg epoch %d, current epoch %d, update local ep", tmq->consumerId,
             head->epoch, epoch);
1547
  }
wmmhello's avatar
wmmhello 已提交
1548
  pParam->pUserFn(tmq, code, pMsg, pParam->pParam);
L
Liu Jicong 已提交
1549

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

dengyihao's avatar
dengyihao 已提交
1552
  taosMemoryFree(pMsg->pEpSet);
L
Liu Jicong 已提交
1553
  taosMemoryFree(pMsg->pData);
1554
  taosMemoryFree(pParam);
L
Liu Jicong 已提交
1555
  return code;
1556 1557
}

L
Liu Jicong 已提交
1558
void tmqBuildConsumeReqImpl(SMqPollReq* pReq, tmq_t* tmq, int64_t timeout, SMqClientTopic* pTopic, SMqClientVg* pVg) {
L
Liu Jicong 已提交
1559 1560 1561 1562
  int32_t groupLen = strlen(tmq->groupId);
  memcpy(pReq->subKey, tmq->groupId, groupLen);
  pReq->subKey[groupLen] = TMQ_SEPARATOR;
  strcpy(pReq->subKey + groupLen + 1, pTopic->topicName);
1563

1564
  pReq->withTbName = tmq->withTbName;
L
Liu Jicong 已提交
1565
  pReq->consumerId = tmq->consumerId;
1566
  pReq->timeout = timeout;
X
Xiaoyu Wang 已提交
1567
  pReq->epoch = tmq->epoch;
L
Liu Jicong 已提交
1568
  /*pReq->currentOffset = reqOffset;*/
H
Haojun Liao 已提交
1569
  pReq->reqOffset = pVg->offsetInfo.currentOffset;
D
dapan1121 已提交
1570
  pReq->head.vgId = pVg->vgId;
1571 1572
  pReq->useSnapshot = tmq->useSnapshot;
  pReq->reqId = generateRequestId();
1573 1574
}

L
Liu Jicong 已提交
1575 1576
SMqMetaRspObj* tmqBuildMetaRspFromWrapper(SMqPollRspWrapper* pWrapper) {
  SMqMetaRspObj* pRspObj = taosMemoryCalloc(1, sizeof(SMqMetaRspObj));
L
Liu Jicong 已提交
1577
  pRspObj->resType = RES_TYPE__TMQ_META;
L
Liu Jicong 已提交
1578 1579 1580 1581 1582 1583 1584 1585
  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;
}

1586
SMqRspObj* tmqBuildRspFromWrapper(SMqPollRspWrapper* pWrapper, SMqClientVg* pVg, int64_t* numOfRows) {
L
Liu Jicong 已提交
1587 1588
  SMqRspObj* pRspObj = taosMemoryCalloc(1, sizeof(SMqRspObj));
  pRspObj->resType = RES_TYPE__TMQ;
1589

1590
  (*numOfRows) = 0;
L
Liu Jicong 已提交
1591 1592
  tstrncpy(pRspObj->topic, pWrapper->topicHandle->topicName, TSDB_TOPIC_FNAME_LEN);
  tstrncpy(pRspObj->db, pWrapper->topicHandle->db, TSDB_DB_FNAME_LEN);
1593

L
Liu Jicong 已提交
1594
  pRspObj->vgId = pWrapper->vgHandle->vgId;
L
Liu Jicong 已提交
1595
  pRspObj->resIter = -1;
L
Liu Jicong 已提交
1596
  memcpy(&pRspObj->rsp, &pWrapper->dataRsp, sizeof(SMqDataRsp));
L
Liu Jicong 已提交
1597

L
Liu Jicong 已提交
1598 1599
  pRspObj->resInfo.totalRows = 0;
  pRspObj->resInfo.precision = TSDB_TIME_PRECISION_MILLI;
1600

L
Liu Jicong 已提交
1601
  if (!pWrapper->dataRsp.withSchema) {
L
Liu Jicong 已提交
1602 1603
    setResSchemaInfo(&pRspObj->resInfo, pWrapper->topicHandle->schema.pSchema, pWrapper->topicHandle->schema.nCols);
  }
L
Liu Jicong 已提交
1604

1605
  // extract the rows in this data packet
X
Xiaoyu Wang 已提交
1606
  for (int32_t i = 0; i < pRspObj->rsp.blockNum; ++i) {
1607
    SRetrieveTableRsp* pRetrieve = (SRetrieveTableRsp*)taosArrayGetP(pRspObj->rsp.blockData, i);
X
Xiaoyu Wang 已提交
1608
    int64_t            rows = htobe64(pRetrieve->numOfRows);
1609
    pVg->numOfRows += rows;
1610
    (*numOfRows) += rows;
1611 1612
  }

L
Liu Jicong 已提交
1613
  return pRspObj;
X
Xiaoyu Wang 已提交
1614 1615
}

L
Liu Jicong 已提交
1616 1617
SMqTaosxRspObj* tmqBuildTaosxRspFromWrapper(SMqPollRspWrapper* pWrapper) {
  SMqTaosxRspObj* pRspObj = taosMemoryCalloc(1, sizeof(SMqTaosxRspObj));
1618
  pRspObj->resType = RES_TYPE__TMQ_METADATA;
L
Liu Jicong 已提交
1619 1620 1621 1622
  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;
1623
  memcpy(&pRspObj->rsp, &pWrapper->taosxRsp, sizeof(STaosxRsp));
L
Liu Jicong 已提交
1624 1625 1626

  pRspObj->resInfo.totalRows = 0;
  pRspObj->resInfo.precision = TSDB_TIME_PRECISION_MILLI;
1627
  if (!pWrapper->taosxRsp.withSchema) {
L
Liu Jicong 已提交
1628 1629 1630 1631 1632 1633
    setResSchemaInfo(&pRspObj->resInfo, pWrapper->topicHandle->schema.pSchema, pWrapper->topicHandle->schema.nCols);
  }

  return pRspObj;
}

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

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

  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];
H
Haojun Liao 已提交
1693
  tFormatOffset(offsetFormatBuf, tListLen(offsetFormatBuf), &pVg->offsetInfo.currentOffset);
1694

X
Xiaoyu Wang 已提交
1695 1696
  tscDebug("consumer:0x%" PRIx64 " send poll to %s vgId:%d, epoch %d, req:%s, reqId:0x%" PRIx64, pTmq->consumerId,
           pTopic->topicName, pVg->vgId, pTmq->epoch, offsetFormatBuf, req.reqId);
1697 1698 1699 1700 1701 1702 1703 1704
  asyncSendMsgToServer(pTmq->pTscObj->pAppInfo->pTransporter, &pVg->epSet, &transporterId, sendInfo);

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

  return TSDB_CODE_SUCCESS;
}

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

  for (int i = 0; i < numOfTopics; i++) {
X
Xiaoyu Wang 已提交
1711
    SMqClientTopic* pTopic = taosArrayGet(tmq->clientTopics, i);
X
Xiaoyu Wang 已提交
1712
    int32_t         numOfVg = taosArrayGetSize(pTopic->vgs);
1713 1714

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

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

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

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

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

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

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

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

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

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

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

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

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

1809
        // update the local offset value only for the returned values.
H
Haojun Liao 已提交
1810
        pVg->offsetInfo.currentOffset = pDataRsp->rspOffset;
X
Xiaoyu Wang 已提交
1811
        atomic_store_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE);
H
Haojun Liao 已提交
1812

1813 1814 1815
        char buf[80];
        tFormatOffset(buf, 80, &pDataRsp->rspOffset);
        if (pDataRsp->blockNum == 0) {
X
Xiaoyu Wang 已提交
1816 1817 1818
          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);
1819
          pRspWrapper = tmqFreeRspWrapper(pRspWrapper);
L
Liu Jicong 已提交
1820
          taosFreeQitem(pollRspWrapper);
1821
        } else {  // build rsp
X
Xiaoyu Wang 已提交
1822
          int64_t    numOfRows = 0;
1823
          SMqRspObj* pRsp = tmqBuildRspFromWrapper(pollRspWrapper, pVg, &numOfRows);
H
Haojun Liao 已提交
1824 1825
          tmq->totalRows += numOfRows;

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

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

L
Liu Jicong 已提交
1845
      if (pollRspWrapper->metaRsp.head.epoch == consumerEpoch) {
L
Liu Jicong 已提交
1846
        SMqClientVg* pVg = pollRspWrapper->vgHandle;
H
Haojun Liao 已提交
1847
        pVg->offsetInfo.currentOffset = pollRspWrapper->metaRsp.rspOffset;
L
Liu Jicong 已提交
1848 1849
        atomic_store_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE);
        // build rsp
L
Liu Jicong 已提交
1850
        SMqMetaRspObj* pRsp = tmqBuildMetaRspFromWrapper(pollRspWrapper);
L
Liu Jicong 已提交
1851 1852 1853
        taosFreeQitem(pollRspWrapper);
        return pRsp;
      } else {
H
Haojun Liao 已提交
1854
        tscDebug("consumer:0x%" PRIx64 " vgId:%d msg discard since epoch mismatch: msg epoch %d, consumer epoch %d",
1855
                 tmq->consumerId, pollRspWrapper->vgId, pollRspWrapper->metaRsp.head.epoch, consumerEpoch);
1856
        pRspWrapper = tmqFreeRspWrapper(pRspWrapper);
L
Liu Jicong 已提交
1857
        taosFreeQitem(pollRspWrapper);
X
Xiaoyu Wang 已提交
1858
      }
1859 1860
    } else if (pRspWrapper->tmqRspType == TMQ_MSG_TYPE__TAOSX_RSP) {
      SMqPollRspWrapper* pollRspWrapper = (SMqPollRspWrapper*)pRspWrapper;
X
Xiaoyu Wang 已提交
1861
      int32_t            consumerEpoch = atomic_load_32(&tmq->epoch);
H
Haojun Liao 已提交
1862

L
Liu Jicong 已提交
1863 1864
      if (pollRspWrapper->taosxRsp.head.epoch == consumerEpoch) {
        SMqClientVg* pVg = pollRspWrapper->vgHandle;
H
Haojun Liao 已提交
1865
        pVg->offsetInfo.currentOffset = pollRspWrapper->taosxRsp.rspOffset;
L
Liu Jicong 已提交
1866
        atomic_store_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE);
H
Haojun Liao 已提交
1867

L
Liu Jicong 已提交
1868
        if (pollRspWrapper->taosxRsp.blockNum == 0) {
H
Haojun Liao 已提交
1869 1870
          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 已提交
1871
          pVg->emptyBlockReceiveTs = taosGetTimestampMs();
1872
          pRspWrapper = tmqFreeRspWrapper(pRspWrapper);
H
Haojun Liao 已提交
1873
          taosFreeQitem(pollRspWrapper);
L
Liu Jicong 已提交
1874
          continue;
H
Haojun Liao 已提交
1875
        } else {
X
Xiaoyu Wang 已提交
1876
          pVg->emptyBlockReceiveTs = 0;  // reset the ts
L
Liu Jicong 已提交
1877
        }
wmmhello's avatar
wmmhello 已提交
1878

L
Liu Jicong 已提交
1879
        // build rsp
X
Xiaoyu Wang 已提交
1880
        void*   pRsp = NULL;
1881
        int64_t numOfRows = 0;
L
Liu Jicong 已提交
1882
        if (pollRspWrapper->taosxRsp.createTableNum == 0) {
1883
          pRsp = tmqBuildRspFromWrapper(pollRspWrapper, pVg, &numOfRows);
L
Liu Jicong 已提交
1884
        } else {
wmmhello's avatar
wmmhello 已提交
1885 1886
          pRsp = tmqBuildTaosxRspFromWrapper(pollRspWrapper);
        }
H
Haojun Liao 已提交
1887

1888 1889
        tmq->totalRows += numOfRows;

H
Haojun Liao 已提交
1890
        char buf[80];
H
Haojun Liao 已提交
1891
        tFormatOffset(buf, 80, &pVg->offsetInfo.currentOffset);
H
Haojun Liao 已提交
1892
        tscDebug("consumer:0x%" PRIx64 " process taosx poll rsp, vgId:%d, offset:%s, blocks:%d, rows:%" PRId64
X
Xiaoyu Wang 已提交
1893
                 ", vg total:%" PRId64 " total:%" PRId64 " reqId:0x%" PRIx64,
H
Haojun Liao 已提交
1894
                 tmq->consumerId, pVg->vgId, buf, pollRspWrapper->dataRsp.blockNum, numOfRows, pVg->numOfRows,
H
Haojun Liao 已提交
1895
                 tmq->totalRows, pollRspWrapper->reqId);
H
Haojun Liao 已提交
1896 1897

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

L
Liu Jicong 已提交
1900
      } else {
H
Haojun Liao 已提交
1901
        tscDebug("consumer:0x%" PRIx64 " vgId:%d msg discard since epoch mismatch: msg epoch %d, consumer epoch %d",
1902
                 tmq->consumerId, pollRspWrapper->vgId, pollRspWrapper->taosxRsp.head.epoch, consumerEpoch);
1903
        pRspWrapper = tmqFreeRspWrapper(pRspWrapper);
L
Liu Jicong 已提交
1904 1905
        taosFreeQitem(pollRspWrapper);
      }
X
Xiaoyu Wang 已提交
1906
    } else {
H
Haojun Liao 已提交
1907 1908
      tscDebug("consumer:0x%" PRIx64 " not data msg received", tmq->consumerId);

X
Xiaoyu Wang 已提交
1909
      bool reset = false;
1910 1911
      tmqHandleNoPollRsp(tmq, pRspWrapper, &reset);
      taosFreeQitem(pRspWrapper);
X
Xiaoyu Wang 已提交
1912
      if (pollIfReset && reset) {
1913
        tscDebug("consumer:0x%" PRIx64 ", reset and repoll", tmq->consumerId);
1914
        tmqPollImpl(tmq, timeout);
X
Xiaoyu Wang 已提交
1915 1916 1917 1918 1919
      }
    }
  }
}

1920
TAOS_RES* tmq_consumer_poll(tmq_t* tmq, int64_t timeout) {
L
Liu Jicong 已提交
1921 1922
  void*   rspObj;
  int64_t startTime = taosGetTimestampMs();
L
Liu Jicong 已提交
1923

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

1927 1928 1929
#if 0
  tmqHandleAllDelayedTask(tmq);
  tmqPollImpl(tmq, timeout);
1930
  rspObj = tmqHandleAllRsp(tmq, timeout, false);
L
Liu Jicong 已提交
1931 1932
  if (rspObj) {
    return (TAOS_RES*)rspObj;
L
fix  
Liu Jicong 已提交
1933
  }
1934
#endif
X
Xiaoyu Wang 已提交
1935

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

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

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

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

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

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

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

1987 1988 1989 1990 1991 1992 1993 1994 1995 1996 1997 1998 1999 2000
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);
2001
    }
2002
  }
2003

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

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

2011 2012 2013 2014 2015 2016
  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;
2017 2018 2019
      }
    }

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

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

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

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

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

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

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

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

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

2127 2128 2129 2130 2131 2132
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 已提交
2133 2134
}

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

2141 2142 2143 2144 2145 2146 2147 2148 2149
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 已提交
2150 2151
  } else {
    asyncCommitOffset(tmq, pRes, commitCallBackFn, pInfo);
2152 2153
  }

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

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

X
Xiaoyu Wang 已提交
2160
  tscDebug("consumer:0x%" PRIx64 " sync commit done, code:%s", tmq->consumerId, tstrerror(code));
2161 2162 2163 2164 2165 2166 2167 2168 2169 2170 2171 2172
  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);
2173
    doUpdateLocalEp(pTmq, head->epoch, &rsp);
2174 2175 2176
    tDeleteSMqAskEpRsp(&rsp);
  }

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

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 已提交
2204
  tsem_init(&pInfo->sem, 0, 0);
2205 2206

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

  return NULL;
H
Haojun Liao 已提交
2335 2336 2337 2338 2339 2340 2341 2342 2343 2344 2345 2346 2347 2348 2349 2350 2351 2352 2353 2354 2355 2356 2357 2358 2359 2360 2361 2362 2363 2364 2365 2366 2367 2368 2369 2370 2371 2372 2373 2374 2375 2376 2377 2378 2379 2380 2381 2382 2383 2384 2385 2386 2387 2388 2389 2390 2391 2392 2393 2394 2395 2396 2397 2398 2399 2400 2401 2402 2403 2404 2405 2406 2407 2408 2409 2410 2411 2412 2413 2414 2415 2416 2417 2418 2419 2420 2421 2422 2423 2424 2425 2426 2427 2428 2429 2430 2431 2432 2433 2434 2435 2436 2437 2438 2439 2440 2441 2442 2443 2444 2445 2446 2447 2448 2449 2450 2451 2452 2453 2454 2455 2456 2457 2458 2459 2460 2461 2462 2463 2464 2465 2466 2467 2468 2469 2470 2471 2472 2473 2474 2475 2476 2477 2478 2479 2480
}

static SMqClientTopic* getTopicByName(tmq_t* tmq, const char* pTopicName) {
  int32_t numOfTopics = taosArrayGetSize(tmq->clientTopics);
  for (int32_t i = 0; i < numOfTopics; ++i) {
    SMqClientTopic* pTopic = taosArrayGet(tmq->clientTopics, i);
    if (strcmp(pTopic->topicName, pTopicName) != 0) {
      continue;
    }

    return pTopic;
  }

  tscError("consumer:0x%" PRIx64 ", failed to find topic:%s", tmq->consumerId, pTopicName);
  return NULL;
}

int32_t tmq_get_topic_assignment(tmq_t* tmq, const char* pTopicName, tmq_topic_assignment** assignment,
                              int32_t* numOfAssignment) {
  *numOfAssignment = 0;
  *assignment = NULL;

  SMqClientTopic* pTopic = getTopicByName(tmq, pTopicName);
  if (pTopic == NULL) {
    return TSDB_CODE_INVALID_PARA;
  }

  // in case of snapshot is opened, no valid offset will return
  *numOfAssignment = taosArrayGetSize(pTopic->vgs);
  for (int32_t j = 0; j < (*numOfAssignment); ++j) {
    SMqClientVg* pClientVg = taosArrayGet(pTopic->vgs, j);

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

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

  return TSDB_CODE_SUCCESS;
}

int32_t tmq_offset_seek(tmq_t* tmq, const char* pTopicName, int32_t vgroupHandle, int64_t offset) {
  if (tmq == NULL) {
    return TSDB_CODE_INVALID_PARA;
  }

  SMqClientTopic* pTopic = getTopicByName(tmq, pTopicName);
  if (pTopic == NULL) {
    return TSDB_CODE_INVALID_PARA;
  }

  SMqClientVg* pVg = NULL;

  int32_t numOfVgs = taosArrayGetSize(pTopic->vgs);
  for(int32_t i= 0; i < numOfVgs; ++i) {
    SMqClientVg* pClientVg = taosArrayGet(pTopic->vgs, i);
    if (pClientVg->vgId == vgroupHandle) {
      pVg = pClientVg;
      break;
    }
  }

  if (pVg == NULL) {
    return TSDB_CODE_INVALID_PARA;
  }

  if (offset < pVg->offsetInfo.walVerBegin|| offset > pVg->offsetInfo.walVerEnd) {
    return TSDB_CODE_INVALID_PARA;
  }

  return 0;
#if 0
  //  tmq_commit_sync(tmq, );
  {
    SMqCommitCbParamSet* pParamSet = taosMemoryCalloc(1, sizeof(SMqCommitCbParamSet));
    if (pParamSet == NULL) {
//      pCommitFp(tmq, TSDB_CODE_OUT_OF_MEMORY, userParam);
      return -1;
    }

    pParamSet->refId = tmq->refId;
    pParamSet->epoch = tmq->epoch;
    pParamSet->callbackFn = pCommitFp;
    pParamSet->userParam = userParam;

    int32_t numOfTopics = taosArrayGetSize(tmq->clientTopics);

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

    int32_t i = 0;
    for (; i < numOfTopics; i++) {
      SMqClientTopic* pTopic = taosArrayGet(tmq->clientTopics, i);
      if (strcmp(pTopic->topicName, pTopicName) == 0) {
        break;
      }
    }

    if (i == numOfTopics) {
      tscWarn("consumer:0x%" PRIx64 " failed to find the specified topic:%s, total topics:%d", tmq->consumerId,
              pTopicName, numOfTopics);
      taosMemoryFree(pParamSet);
      pCommitFp(tmq, TSDB_CODE_SUCCESS, userParam);
      return;
    }

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

    if (j == numOfVgroups) {
      tscWarn("consumer:0x%" PRIx64 " failed to find the specified vgId:%d, total Vgs:%d, topic:%s", tmq->consumerId,
              vgId, numOfVgroups, pTopicName);
      taosMemoryFree(pParamSet);
      pCommitFp(tmq, TSDB_CODE_SUCCESS, userParam);
      return;
    }

    SMqClientVg* pVg = taosArrayGet(pTopic->vgs, j);
    if (pVg->offsetInfo.currentOffset.type > 0 && !tOffsetEqual(&pVg->offsetInfo.currentOffset, &pVg->offsetInfo.committedOffset)) {
      code = doSendCommitMsg(tmq, pVg, pTopic->topicName, pParamSet, j, numOfVgroups);

      // 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);
    }
  }
#endif

2481
}