clientTmq.c 77.0 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

  // status
wmmhello's avatar
wmmhello 已提交
90
  SRWLatch        lock;
L
Liu Jicong 已提交
91 92
  int8_t  status;
  int32_t epoch;
L
Liu Jicong 已提交
93 94
#if 0
  int8_t  epStatus;
L
Liu Jicong 已提交
95
  int32_t epSkipCnt;
L
Liu Jicong 已提交
96
#endif
97
  // poll info
X
Xiaoyu Wang 已提交
98 99
  int64_t pollCnt;
  int64_t totalRows;
L
Liu Jicong 已提交
100

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

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

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

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

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

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

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

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

L
Liu Jicong 已提交
172
typedef struct {
wmmhello's avatar
wmmhello 已提交
173 174
//  int64_t refId;
//  int32_t epoch;
L
Liu Jicong 已提交
175 176
  tsem_t  rspSem;
  int32_t rspErr;
L
Liu Jicong 已提交
177
} SMqSubscribeCbParam;
L
Liu Jicong 已提交
178

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

L
Liu Jicong 已提交
186
typedef struct {
187 188
  int64_t         refId;
  int32_t         epoch;
wmmhello's avatar
wmmhello 已提交
189 190 191
  char            topicName[TSDB_TOPIC_FNAME_LEN];
//  SMqClientVg*    pVg;
//  SMqClientTopic* pTopic;
L
Liu Jicong 已提交
192
  int32_t         vgId;
X
Xiaoyu Wang 已提交
193
  uint64_t        requestId;  // request id for debug purpose
X
Xiaoyu Wang 已提交
194
} SMqPollCbParam;
195

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

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

216 217 218 219 220
typedef struct SSyncCommitInfo {
  tsem_t  sem;
  int32_t code;
} SSyncCommitInfo;

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

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

237
  conf->withTbName = false;
L
Liu Jicong 已提交
238
  conf->autoCommit = true;
239
  conf->autoCommitInterval = DEFAULT_AUTO_COMMIT_INTERVAL;
240
  conf->resetOffset = TMQ_OFFSET__RESET_EARLIEAST;
241
  conf->hbBgEnable = true;
242

243 244 245
  return conf;
}

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

L
Liu Jicong 已提交
370
  return TMQ_CONF_UNKNOWN;
371 372
}

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

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

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

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

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

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

  return NULL;
L
Liu Jicong 已提交
421
}
422

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

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

L
Liu Jicong 已提交
461
  taosMemoryFree(pParam->pOffset);
L
Liu Jicong 已提交
462
  taosMemoryFree(pBuf->pData);
dengyihao's avatar
dengyihao 已提交
463
  taosMemoryFree(pBuf->pEpSet);
L
Liu Jicong 已提交
464

465
  commitRspCountDown(pParamSet, pParam->pTmq->consumerId, pParam->topicName, pParam->vgId);
466 467 468
  return 0;
}

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

L
Liu Jicong 已提交
476
  pOffset->val = pVg->currentOffset;
477

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

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

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

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

L
Liu Jicong 已提交
498 499 500 501 502
  void* abuf = POINTER_SHIFT(buf, sizeof(SMsgHead));

  SEncoder encoder;
  tEncoderInit(&encoder, abuf, len);
  tEncodeSTqOffset(&encoder, pOffset);
L
Liu Jicong 已提交
503
  tEncoderClear(&encoder);
L
Liu Jicong 已提交
504 505

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

L
Liu Jicong 已提交
513 514
  pParam->params = pParamSet;
  pParam->pOffset = pOffset;
H
Haojun Liao 已提交
515 516 517
  pParam->vgId = pVg->vgId;
  pParam->pTmq = tmq;

H
Haojun Liao 已提交
518
  tstrncpy(pParam->topicName, pTopicName, tListLen(pParam->topicName));
L
Liu Jicong 已提交
519 520 521 522

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

790
  sendInfo->msgInfo = (SDataBuf){ .pData = pReq, .len = tlen, .handle = NULL };
791 792 793 794 795

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

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

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

OVER:
804
  taosTmrReset(tmqSendHbReq, 1000, param, tmqMgmt.timer, &tmq->hbLiveTimer);
805
  taosReleaseRef(tmqMgmt.rsetId, refId);
806 807
}

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

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

818 819 820 821
  if (qall->numOfItems == 0) {
    taosFreeQall(qall);
    return TSDB_CODE_SUCCESS;
  }
822

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

827
  while (pTaskType != NULL) {
828
    if (*pTaskType == TMQ_DELAYED_TASK__ASK_EP) {
829
      asyncAskEp(pTmq, addToQueueCallbackFn, NULL);
830 831

      int64_t* pRefId = taosMemoryMalloc(sizeof(int64_t));
832
      *pRefId = pTmq->refId;
833

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

      asyncCommitAllOffsets(pTmq, pCallbackFn, pTmq->commitCbUserParam);
840
      int64_t* pRefId = taosMemoryMalloc(sizeof(int64_t));
841
      *pRefId = pTmq->refId;
842

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

L
Liu Jicong 已提交
849
    taosFreeQitem(pTaskType);
850
    taosGetQitem(qall, (void**)&pTaskType);
L
Liu Jicong 已提交
851
  }
852

L
Liu Jicong 已提交
853 854 855 856
  taosFreeQall(qall);
  return 0;
}

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

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

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

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

  return NULL;
L
Liu Jicong 已提交
890 891
}

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

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

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

  taosMemoryFree(pMsg->pEpSet);
L
Liu Jicong 已提交
922 923 924
  tsem_post(&pParam->rspSem);
  return 0;
}
925

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

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

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

961
void tmqFreeImpl(void* handle) {
962 963
  tmq_t*  tmq = (tmq_t*)handle;
  int64_t id = tmq->consumerId;
L
Liu Jicong 已提交
964

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

H
Haojun Liao 已提交
971 972 973 974 975
  if (tmq->delayedTask) {
    taosCloseQueue(tmq->delayedTask);
  }

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

978
  taosArrayDestroyEx(tmq->clientTopics, freeClientVgImpl);
979 980
  taos_close_internal(tmq->pTscObj);
  taosMemoryFree(tmq);
981 982

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

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

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

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

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

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

L
Liu Jicong 已提交
1028 1029
  // init status
  pTmq->status = TMQ_CONSUMER_STATUS__INIT;
L
Liu Jicong 已提交
1030 1031
  pTmq->pollCnt = 0;
  pTmq->epoch = 0;
L
Liu Jicong 已提交
1032

L
Liu Jicong 已提交
1033 1034 1035
  // set conf
  strcpy(pTmq->clientId, conf->clientId);
  strcpy(pTmq->groupId, conf->groupId);
1036
  pTmq->withTbName = conf->withTbName;
L
Liu Jicong 已提交
1037
  pTmq->useSnapshot = conf->snapEnable;
L
Liu Jicong 已提交
1038
  pTmq->autoCommit = conf->autoCommit;
L
Liu Jicong 已提交
1039
  pTmq->autoCommitInterval = conf->autoCommitInterval;
L
Liu Jicong 已提交
1040 1041
  pTmq->commitCb = conf->commitCb;
  pTmq->commitCbUserParam = conf->commitCbUserParam;
L
Liu Jicong 已提交
1042
  pTmq->resetOffsetCfg = conf->resetOffset;
wmmhello's avatar
wmmhello 已提交
1043
  taosInitRWLatch(&pTmq->lock);
L
Liu Jicong 已提交
1044

1045 1046
  pTmq->hbBgEnable = conf->hbBgEnable;

L
Liu Jicong 已提交
1047
  // assign consumerId
L
Liu Jicong 已提交
1048
  pTmq->consumerId = tGenIdPI64();
X
Xiaoyu Wang 已提交
1049

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

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

1065 1066
  pTmq->refId = taosAddRef(tmqMgmt.rsetId, pTmq);
  if (pTmq->refId < 0) {
1067
    goto _failed;
1068 1069
  }

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

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

1084
  return pTmq;
1085

1086 1087
_failed:
  tmqFreeImpl(pTmq);
L
Liu Jicong 已提交
1088
  return NULL;
1089 1090
}

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

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

1102
  req.consumerId = tmq->consumerId;
L
Liu Jicong 已提交
1103
  tstrncpy(req.clientId, tmq->clientId, 256);
L
Liu Jicong 已提交
1104
  tstrncpy(req.cgroup, tmq->groupId, TSDB_CGROUP_LEN);
1105 1106
  req.topicNames = taosArrayInit(sz, sizeof(void*));

1107 1108 1109 1110
  if (req.topicNames == NULL) {
    code = TSDB_CODE_OUT_OF_MEMORY;
    goto FAIL;
  }
1111

L
Liu Jicong 已提交
1112 1113
  for (int32_t i = 0; i < sz; i++) {
    char* topic = taosArrayGetP(container, i);
1114 1115

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

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

    taosArrayPush(req.topicNames, &topicFName);
1126 1127
  }

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

L
Liu Jicong 已提交
1130
  buf = taosMemoryMalloc(tlen);
1131 1132 1133 1134
  if (buf == NULL) {
    code = TSDB_CODE_OUT_OF_MEMORY;
    goto FAIL;
  }
L
Liu Jicong 已提交
1135

1136 1137 1138
  void* abuf = buf;
  tSerializeSCMSubscribeReq(&abuf, &req);

L
Liu Jicong 已提交
1139
  sendInfo = taosMemoryCalloc(1, sizeof(SMsgSendInfo));
1140 1141 1142 1143
  if (sendInfo == NULL) {
    code = TSDB_CODE_OUT_OF_MEMORY;
    goto FAIL;
  }
1144

X
Xiaoyu Wang 已提交
1145
  SMqSubscribeCbParam param = {
L
Liu Jicong 已提交
1146
      .rspErr = 0,
wmmhello's avatar
wmmhello 已提交
1147 1148
//      .refId = tmq->refId,
//      .epoch = tmq->epoch,
X
Xiaoyu Wang 已提交
1149
  };
L
Liu Jicong 已提交
1150

1151
  if (tsem_init(&param.rspSem, 0, 0) != 0) {
wmmhello's avatar
wmmhello 已提交
1152
    code = TSDB_CODE_TSC_INTERNAL_ERROR;
1153 1154
    goto FAIL;
  }
L
Liu Jicong 已提交
1155 1156

  sendInfo->msgInfo = (SDataBuf){
X
Xiaoyu Wang 已提交
1157 1158 1159 1160
      .pData = buf,
      .len = tlen,
      .handle = NULL,
  };
1161

L
Liu Jicong 已提交
1162 1163
  sendInfo->requestId = generateRequestId();
  sendInfo->requestObjRefId = 0;
L
Liu Jicong 已提交
1164 1165
  sendInfo->param = &param;
  sendInfo->fp = tmqSubscribeCb;
L
Liu Jicong 已提交
1166
  sendInfo->msgType = TDMT_MND_TMQ_SUBSCRIBE;
L
Liu Jicong 已提交
1167

1168 1169 1170 1171 1172
  SEpSet epSet = getEpSet_s(&tmq->pTscObj->pAppInfo->mgmtEp);

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

L
Liu Jicong 已提交
1173 1174
  // avoid double free if msg is sent
  buf = NULL;
L
Liu Jicong 已提交
1175
  sendInfo = NULL;
L
Liu Jicong 已提交
1176

L
Liu Jicong 已提交
1177 1178
  tsem_wait(&param.rspSem);
  tsem_destroy(&param.rspSem);
1179

1180 1181 1182 1183
  if (param.rspErr != 0) {
    code = param.rspErr;
    goto FAIL;
  }
L
Liu Jicong 已提交
1184

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

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

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

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

L
Liu Jicong 已提交
1211
FAIL:
L
Liu Jicong 已提交
1212
  taosArrayDestroyP(req.topicNames, taosMemoryFree);
L
Liu Jicong 已提交
1213
  taosMemoryFree(buf);
L
Liu Jicong 已提交
1214
  taosMemoryFree(sendInfo);
L
Liu Jicong 已提交
1215

L
Liu Jicong 已提交
1216
  return code;
1217 1218
}

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

wmmhello's avatar
wmmhello 已提交
1224 1225 1226 1227 1228 1229 1230 1231 1232 1233 1234 1235 1236 1237 1238 1239 1240 1241 1242 1243 1244 1245 1246 1247 1248 1249 1250 1251
static SMqClientVg* getVgInfo(tmq_t* tmq, char* topicName, int32_t  vgId){
  int32_t topicNumCur = taosArrayGetSize(tmq->clientTopics);
  for(int i = 0; i < topicNumCur; i++){
    SMqClientTopic* pTopicCur = taosArrayGet(tmq->clientTopics, i);
    if(strcmp(pTopicCur->topicName, topicName) == 0){
      int32_t vgNumCur = taosArrayGetSize(pTopicCur->vgs);
      for (int32_t j = 0; j < vgNumCur; j++) {
        SMqClientVg* pVgCur = taosArrayGet(pTopicCur->vgs, j);
        if(pVgCur->vgId == vgId){
          return pVgCur;
        }
      }
    }
  }
  return NULL;
}

static SMqClientTopic* getTopicInfo(tmq_t* tmq, char* topicName){
  int32_t topicNumCur = taosArrayGetSize(tmq->clientTopics);
  for(int i = 0; i < topicNumCur; i++){
    SMqClientTopic* pTopicCur = taosArrayGet(tmq->clientTopics, i);
    if(strcmp(pTopicCur->topicName, topicName) == 0){
      return pTopicCur;
    }
  }
  return NULL;
}

D
dapan1121 已提交
1252
int32_t tmqPollCb(void* param, SDataBuf* pMsg, int32_t code) {
X
Xiaoyu Wang 已提交
1253
  SMqPollCbParam* pParam = (SMqPollCbParam*)param;
1254 1255

  int64_t         refId = pParam->refId;
wmmhello's avatar
wmmhello 已提交
1256 1257
//  SMqClientVg*    pVg = pParam->pVg;
//  SMqClientTopic* pTopic = pParam->pTopic;
1258

1259
  tmq_t* tmq = taosAcquireRef(tmqMgmt.rsetId, refId);
1260 1261 1262
  if (tmq == NULL) {
    taosMemoryFree(pParam);
    taosMemoryFree(pMsg->pData);
dengyihao's avatar
dengyihao 已提交
1263
    taosMemoryFree(pMsg->pEpSet);
1264 1265 1266 1267
    terrno = TSDB_CODE_TMQ_CONSUMER_CLOSED;
    return -1;
  }

H
Haojun Liao 已提交
1268 1269 1270 1271
  int32_t  epoch = pParam->epoch;
  int32_t  vgId = pParam->vgId;
  uint64_t requestId = pParam->requestId;

L
Liu Jicong 已提交
1272
  if (code != 0) {
L
Liu Jicong 已提交
1273
    if (pMsg->pData) taosMemoryFree(pMsg->pData);
dengyihao's avatar
dengyihao 已提交
1274 1275
    if (pMsg->pEpSet) taosMemoryFree(pMsg->pEpSet);

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

L
Liu Jicong 已提交
1290 1291
      pRspWrapper->tmqRspType = TMQ_MSG_TYPE__END_RSP;
      taosWriteQitem(tmq->mqueue, pRspWrapper);
X
Xiaoyu Wang 已提交
1292
    } else if (code == TSDB_CODE_WAL_LOG_NOT_EXIST) {  // poll data while insert
1293
      taosMsleep(500);
wmmhello's avatar
wmmhello 已提交
1294 1295 1296
    } else{
      tscError("consumer:0x%" PRIx64 " msg from vgId:%d discarded, epoch %d, since %s, reqId:0x%" PRIx64, tmq->consumerId,
               vgId, epoch, tstrerror(code), requestId);
L
Liu Jicong 已提交
1297
    }
H
Haojun Liao 已提交
1298

L
fix txn  
Liu Jicong 已提交
1299
    goto CREATE_MSG_FAIL;
L
Liu Jicong 已提交
1300 1301
  }

X
Xiaoyu Wang 已提交
1302 1303 1304
  int32_t msgEpoch = ((SMqRspHead*)pMsg->pData)->epoch;
  int32_t tmqEpoch = atomic_load_32(&tmq->epoch);
  if (msgEpoch < tmqEpoch) {
L
Liu Jicong 已提交
1305
    // do not write into queue since updating epoch reset
X
Xiaoyu Wang 已提交
1306 1307
    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 已提交
1308 1309
            tmq->consumerId, vgId, msgEpoch, tmqEpoch, requestId);

1310
    tsem_post(&tmq->rspSem);
1311 1312
    taosReleaseRef(tmqMgmt.rsetId, refId);

L
Liu Jicong 已提交
1313
    taosMemoryFree(pMsg->pData);
dengyihao's avatar
dengyihao 已提交
1314
    taosMemoryFree(pMsg->pEpSet);
wmmhello's avatar
wmmhello 已提交
1315 1316
    taosMemoryFree(pParam);

X
Xiaoyu Wang 已提交
1317 1318 1319 1320
    return 0;
  }

  if (msgEpoch != tmqEpoch) {
H
Haojun Liao 已提交
1321 1322
    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 已提交
1323 1324
  }

L
Liu Jicong 已提交
1325 1326 1327
  // handle meta rsp
  int8_t rspType = ((SMqRspHead*)pMsg->pData)->mqMsgType;

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

L
Liu Jicong 已提交
1337
  pRspWrapper->tmqRspType = rspType;
wmmhello's avatar
wmmhello 已提交
1338 1339
//  pRspWrapper->vgHandle = pVg;
//  pRspWrapper->topicHandle = pTopic;
H
Haojun Liao 已提交
1340
  pRspWrapper->reqId = requestId;
1341
  pRspWrapper->pEpset = pMsg->pEpSet;
wmmhello's avatar
wmmhello 已提交
1342 1343
  pRspWrapper->vgId = vgId;
  strcpy(pRspWrapper->topicName, pParam->topicName);
L
Liu Jicong 已提交
1344

1345
  pMsg->pEpSet = NULL;
L
Liu Jicong 已提交
1346
  if (rspType == TMQ_MSG_TYPE__POLL_RSP) {
L
Liu Jicong 已提交
1347 1348 1349
    SDecoder decoder;
    tDecoderInit(&decoder, POINTER_SHIFT(pMsg->pData, sizeof(SMqRspHead)), pMsg->len - sizeof(SMqRspHead));
    tDecodeSMqDataRsp(&decoder, &pRspWrapper->dataRsp);
wmmhello's avatar
wmmhello 已提交
1350
    tDecoderClear(&decoder);
L
Liu Jicong 已提交
1351
    memcpy(&pRspWrapper->dataRsp, pMsg->pData, sizeof(SMqRspHead));
L
Liu Jicong 已提交
1352

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

L
Liu Jicong 已提交
1373
  taosMemoryFree(pMsg->pData);
H
Haojun Liao 已提交
1374
  taosWriteQitem(tmq->mqueue, pRspWrapper);
L
Liu Jicong 已提交
1375

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

1380
  tsem_post(&tmq->rspSem);
1381
  taosReleaseRef(tmqMgmt.rsetId, refId);
wmmhello's avatar
wmmhello 已提交
1382
  taosMemoryFree(pParam);
1383

L
Liu Jicong 已提交
1384
  return 0;
H
Haojun Liao 已提交
1385

L
fix txn  
Liu Jicong 已提交
1386
CREATE_MSG_FAIL:
L
Liu Jicong 已提交
1387
  if (epoch == tmq->epoch) {
wmmhello's avatar
wmmhello 已提交
1388 1389 1390 1391 1392
    taosWLockLatch(&tmq->lock);
    SMqClientVg* pVg = getVgInfo(tmq, pParam->topicName, vgId);
    if(pVg){
      atomic_store_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE);
    }
wmmhello's avatar
wmmhello 已提交
1393
    taosWUnLockLatch(&tmq->lock);
L
Liu Jicong 已提交
1394
  }
H
Haojun Liao 已提交
1395

1396
  tsem_post(&tmq->rspSem);
1397
  taosReleaseRef(tmqMgmt.rsetId, refId);
wmmhello's avatar
wmmhello 已提交
1398
  taosMemoryFree(pParam);
1399

L
Liu Jicong 已提交
1400
  return -1;
1401 1402
}

H
Haojun Liao 已提交
1403 1404 1405 1406 1407
typedef struct SVgroupSaveInfo {
  STqOffsetVal offset;
  int64_t      numOfRows;
} SVgroupSaveInfo;

H
Haojun Liao 已提交
1408 1409 1410 1411 1412 1413
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 已提交
1414
  char    vgKey[TSDB_TOPIC_FNAME_LEN + 22];
H
Haojun Liao 已提交
1415 1416 1417 1418 1419
  int32_t vgNumGet = taosArrayGetSize(pTopicEp->vgs);

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

wmmhello's avatar
wmmhello 已提交
1420
  tscDebug("consumer:0x%" PRIx64 ", update topic:%s, new numOfVgs:%d", tmq->consumerId, pTopic->topicName, vgNumGet);
H
Haojun Liao 已提交
1421 1422 1423 1424
  pTopic->vgs = taosArrayInit(vgNumGet, sizeof(SMqClientVg));

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

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

X
Xiaoyu Wang 已提交
1429 1430
    int64_t      numOfRows = 0;
    STqOffsetVal offsetNew = {.type = tmq->resetOffsetCfg};
H
Haojun Liao 已提交
1431 1432 1433
    if (pInfo != NULL) {
      offsetNew = pInfo->offset;
      numOfRows = pInfo->numOfRows;
H
Haojun Liao 已提交
1434 1435 1436 1437 1438 1439 1440
    }

    SMqClientVg clientVg = {
        .pollCnt = 0,
        .currentOffset = offsetNew,
        .vgId = pVgEp->vgId,
        .epSet = pVgEp->epSet,
wmmhello's avatar
wmmhello 已提交
1441
        .vgStatus = TMQ_VG_STATUS__IDLE,
H
Haojun Liao 已提交
1442
        .vgSkipCnt = 0,
H
Haojun Liao 已提交
1443
        .emptyBlockReceiveTs = 0,
H
Haojun Liao 已提交
1444
        .numOfRows = numOfRows,
H
Haojun Liao 已提交
1445 1446 1447 1448 1449 1450 1451 1452 1453 1454 1455 1456 1457 1458 1459
    };

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

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

  taosArrayDestroy(pTopic->vgs);
}

1460
static bool doUpdateLocalEp(tmq_t* tmq, int32_t epoch, const SMqAskEpRsp* pRsp) {
1461 1462
  bool set = false;

1463
  int32_t topicNumCur = taosArrayGetSize(tmq->clientTopics);
1464
  int32_t topicNumGet = taosArrayGetSize(pRsp->topics);
1465

X
Xiaoyu Wang 已提交
1466 1467
  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",
1468
           tmq->consumerId, tmq->epoch, epoch, topicNumGet, topicNumCur);
wmmhello's avatar
wmmhello 已提交
1469 1470 1471
  if (epoch <= tmq->epoch) {
    return false;
  }
1472 1473 1474 1475 1476 1477

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

H
Haojun Liao 已提交
1478 1479
  SHashObj* pVgOffsetHashMap = taosHashInit(64, MurmurHash3_32, false, HASH_NO_LOCK);
  if (pVgOffsetHashMap == NULL) {
1480 1481 1482
    taosArrayDestroy(newTopics);
    return false;
  }
1483

H
Haojun Liao 已提交
1484
  // todo extract method
1485 1486 1487 1488 1489
  for (int32_t i = 0; i < topicNumCur; i++) {
    // find old topic
    SMqClientTopic* pTopicCur = taosArrayGet(tmq->clientTopics, i);
    if (pTopicCur->vgs) {
      int32_t vgNumCur = taosArrayGetSize(pTopicCur->vgs);
wmmhello's avatar
wmmhello 已提交
1490
      tscDebug("consumer:0x%" PRIx64 ", current vg num: %d", tmq->consumerId, vgNumCur);
1491 1492
      for (int32_t j = 0; j < vgNumCur; j++) {
        SMqClientVg* pVgCur = taosArrayGet(pTopicCur->vgs, j);
H
Haojun Liao 已提交
1493 1494
        makeTopicVgroupKey(vgKey, pTopicCur->topicName, pVgCur->vgId);

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

wmmhello's avatar
wmmhello 已提交
1500
        SVgroupSaveInfo info = {.offset = pVgCur->currentOffset, .numOfRows = pVgCur->numOfRows};
H
Haojun Liao 已提交
1501
        taosHashPut(pVgOffsetHashMap, vgKey, strlen(vgKey), &info, sizeof(SVgroupSaveInfo));
1502 1503 1504 1505 1506 1507 1508
      }
    }
  }

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

H
Haojun Liao 已提交
1513 1514
  taosHashCleanup(pVgOffsetHashMap);

wmmhello's avatar
wmmhello 已提交
1515
  taosWLockLatch(&tmq->lock);
1516
  // destroy current buffered existed topics info
1517
  if (tmq->clientTopics) {
H
Haojun Liao 已提交
1518
    taosArrayDestroyEx(tmq->clientTopics, freeClientVgInfo);
X
Xiaoyu Wang 已提交
1519
  }
H
Haojun Liao 已提交
1520
  tmq->clientTopics = newTopics;
wmmhello's avatar
wmmhello 已提交
1521
  taosWUnLockLatch(&tmq->lock);
1522

X
Xiaoyu Wang 已提交
1523
  int8_t flag = (topicNumGet == 0) ? TMQ_CONSUMER_STATUS__NO_TOPIC : TMQ_CONSUMER_STATUS__READY;
H
Haojun Liao 已提交
1524
  atomic_store_8(&tmq->status, flag);
X
Xiaoyu Wang 已提交
1525
  atomic_store_32(&tmq->epoch, epoch);
H
Haojun Liao 已提交
1526

1527
  tscDebug("consumer:0x%" PRIx64 " update topic info completed", tmq->consumerId);
X
Xiaoyu Wang 已提交
1528 1529 1530
  return set;
}

1531
int32_t askEpCallbackFn(void* param, SDataBuf* pMsg, int32_t code) {
1532
  SMqAskEpCbParam* pParam = (SMqAskEpCbParam*)param;
1533 1534 1535
  tmq_t*           tmq = taosAcquireRef(tmqMgmt.rsetId, pParam->refId);

  if (tmq == NULL) {
1536
    terrno = TSDB_CODE_TMQ_CONSUMER_CLOSED;
wmmhello's avatar
wmmhello 已提交
1537
//    pParam->pUserFn(tmq, terrno, NULL, pParam->pParam);
1538

1539
    taosMemoryFree(pMsg->pData);
dengyihao's avatar
dengyihao 已提交
1540
    taosMemoryFree(pMsg->pEpSet);
1541 1542
    taosMemoryFree(pParam);
    return terrno;
1543 1544
  }

H
Haojun Liao 已提交
1545
  if (code != TSDB_CODE_SUCCESS) {
1546 1547 1548 1549 1550 1551 1552 1553 1554
    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;
1555
  }
L
Liu Jicong 已提交
1556

L
Liu Jicong 已提交
1557
  // tmq's epoch is monotonically increase,
L
Liu Jicong 已提交
1558
  // so it's safe to discard any old epoch msg.
L
Liu Jicong 已提交
1559
  // Epoch will only increase when received newer epoch ep msg
L
Liu Jicong 已提交
1560 1561 1562
  SMqRspHead* head = pMsg->pData;
  int32_t     epoch = atomic_load_32(&tmq->epoch);
  if (head->epoch <= epoch) {
1563 1564
    tscDebug("consumer:0x%" PRIx64 ", recv ep, msg epoch %d, current epoch %d, no need to update local ep",
             tmq->consumerId, head->epoch, epoch);
1565

1566 1567 1568 1569 1570 1571 1572 1573
    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 已提交
1574
  } else {
1575 1576
    tscDebug("consumer:0x%" PRIx64 ", recv ep, msg epoch %d, current epoch %d, update local ep", tmq->consumerId,
             head->epoch, epoch);
1577
  }
wmmhello's avatar
wmmhello 已提交
1578
  pParam->pUserFn(tmq, code, pMsg, pParam->pParam);
L
Liu Jicong 已提交
1579

1580 1581
  taosReleaseRef(tmqMgmt.rsetId, pParam->refId);

dengyihao's avatar
dengyihao 已提交
1582
  taosMemoryFree(pMsg->pEpSet);
L
Liu Jicong 已提交
1583
  taosMemoryFree(pMsg->pData);
1584
  taosMemoryFree(pParam);
L
Liu Jicong 已提交
1585
  return code;
1586 1587
}

L
Liu Jicong 已提交
1588
void tmqBuildConsumeReqImpl(SMqPollReq* pReq, tmq_t* tmq, int64_t timeout, SMqClientTopic* pTopic, SMqClientVg* pVg) {
L
Liu Jicong 已提交
1589 1590 1591 1592
  int32_t groupLen = strlen(tmq->groupId);
  memcpy(pReq->subKey, tmq->groupId, groupLen);
  pReq->subKey[groupLen] = TMQ_SEPARATOR;
  strcpy(pReq->subKey + groupLen + 1, pTopic->topicName);
1593

1594
  pReq->withTbName = tmq->withTbName;
L
Liu Jicong 已提交
1595
  pReq->consumerId = tmq->consumerId;
1596
  pReq->timeout = timeout;
X
Xiaoyu Wang 已提交
1597
  pReq->epoch = tmq->epoch;
L
Liu Jicong 已提交
1598
  /*pReq->currentOffset = reqOffset;*/
L
Liu Jicong 已提交
1599
  pReq->reqOffset = pVg->currentOffset;
D
dapan1121 已提交
1600
  pReq->head.vgId = pVg->vgId;
1601 1602
  pReq->useSnapshot = tmq->useSnapshot;
  pReq->reqId = generateRequestId();
1603 1604
}

L
Liu Jicong 已提交
1605 1606
SMqMetaRspObj* tmqBuildMetaRspFromWrapper(SMqPollRspWrapper* pWrapper) {
  SMqMetaRspObj* pRspObj = taosMemoryCalloc(1, sizeof(SMqMetaRspObj));
L
Liu Jicong 已提交
1607
  pRspObj->resType = RES_TYPE__TMQ_META;
L
Liu Jicong 已提交
1608 1609 1610 1611 1612 1613 1614 1615
  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;
}

1616
SMqRspObj* tmqBuildRspFromWrapper(SMqPollRspWrapper* pWrapper, SMqClientVg* pVg, int64_t* numOfRows) {
L
Liu Jicong 已提交
1617 1618
  SMqRspObj* pRspObj = taosMemoryCalloc(1, sizeof(SMqRspObj));
  pRspObj->resType = RES_TYPE__TMQ;
1619

1620
  (*numOfRows) = 0;
L
Liu Jicong 已提交
1621 1622
  tstrncpy(pRspObj->topic, pWrapper->topicHandle->topicName, TSDB_TOPIC_FNAME_LEN);
  tstrncpy(pRspObj->db, pWrapper->topicHandle->db, TSDB_DB_FNAME_LEN);
1623

L
Liu Jicong 已提交
1624
  pRspObj->vgId = pWrapper->vgHandle->vgId;
L
Liu Jicong 已提交
1625
  pRspObj->resIter = -1;
L
Liu Jicong 已提交
1626
  memcpy(&pRspObj->rsp, &pWrapper->dataRsp, sizeof(SMqDataRsp));
L
Liu Jicong 已提交
1627

L
Liu Jicong 已提交
1628 1629
  pRspObj->resInfo.totalRows = 0;
  pRspObj->resInfo.precision = TSDB_TIME_PRECISION_MILLI;
1630

L
Liu Jicong 已提交
1631
  if (!pWrapper->dataRsp.withSchema) {
L
Liu Jicong 已提交
1632 1633
    setResSchemaInfo(&pRspObj->resInfo, pWrapper->topicHandle->schema.pSchema, pWrapper->topicHandle->schema.nCols);
  }
L
Liu Jicong 已提交
1634

1635
  // extract the rows in this data packet
X
Xiaoyu Wang 已提交
1636
  for (int32_t i = 0; i < pRspObj->rsp.blockNum; ++i) {
1637
    SRetrieveTableRsp* pRetrieve = (SRetrieveTableRsp*)taosArrayGetP(pRspObj->rsp.blockData, i);
X
Xiaoyu Wang 已提交
1638
    int64_t            rows = htobe64(pRetrieve->numOfRows);
1639
    pVg->numOfRows += rows;
1640
    (*numOfRows) += rows;
1641 1642
  }

L
Liu Jicong 已提交
1643
  return pRspObj;
X
Xiaoyu Wang 已提交
1644 1645
}

L
Liu Jicong 已提交
1646 1647
SMqTaosxRspObj* tmqBuildTaosxRspFromWrapper(SMqPollRspWrapper* pWrapper) {
  SMqTaosxRspObj* pRspObj = taosMemoryCalloc(1, sizeof(SMqTaosxRspObj));
1648
  pRspObj->resType = RES_TYPE__TMQ_METADATA;
L
Liu Jicong 已提交
1649 1650 1651 1652
  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;
1653
  memcpy(&pRspObj->rsp, &pWrapper->taosxRsp, sizeof(STaosxRsp));
L
Liu Jicong 已提交
1654 1655 1656

  pRspObj->resInfo.totalRows = 0;
  pRspObj->resInfo.precision = TSDB_TIME_PRECISION_MILLI;
1657
  if (!pWrapper->taosxRsp.withSchema) {
L
Liu Jicong 已提交
1658 1659 1660 1661 1662 1663
    setResSchemaInfo(&pRspObj->resInfo, pWrapper->topicHandle->schema.pSchema, pWrapper->topicHandle->schema.nCols);
  }

  return pRspObj;
}

1664 1665 1666 1667 1668 1669 1670 1671 1672 1673 1674 1675 1676 1677 1678 1679 1680 1681 1682 1683 1684 1685 1686 1687 1688 1689 1690 1691 1692 1693 1694 1695 1696
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;
wmmhello's avatar
wmmhello 已提交
1697 1698 1699
//  pParam->pVg = pVg;  // pVg may be released,fix it
//  pParam->pTopic = pTopic;
  strcpy(pParam->topicName, pTopic->topicName);
1700
  pParam->vgId = pVg->vgId;
H
Haojun Liao 已提交
1701
  pParam->requestId = req.reqId;
1702 1703 1704 1705 1706 1707 1708 1709

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

H
Haojun Liao 已提交
1710
  sendInfo->msgInfo = (SDataBuf){ .pData = msg, .len = msgSize, .handle = NULL };
1711 1712 1713 1714 1715 1716 1717 1718 1719 1720 1721

  sendInfo->requestId = req.reqId;
  sendInfo->requestObjRefId = 0;
  sendInfo->param = pParam;
  sendInfo->fp = tmqPollCb;
  sendInfo->msgType = TDMT_VND_TMQ_CONSUME;

  int64_t transporterId = 0;
  char    offsetFormatBuf[80];
  tFormatOffset(offsetFormatBuf, tListLen(offsetFormatBuf), &pVg->currentOffset);

X
Xiaoyu Wang 已提交
1722 1723
  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);
1724 1725 1726 1727 1728 1729 1730 1731
  asyncSendMsgToServer(pTmq->pTscObj->pAppInfo->pTransporter, &pVg->epSet, &transporterId, sendInfo);

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

  return TSDB_CODE_SUCCESS;
}

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

  for (int i = 0; i < numOfTopics; i++) {
X
Xiaoyu Wang 已提交
1738
    SMqClientTopic* pTopic = taosArrayGet(tmq->clientTopics, i);
X
Xiaoyu Wang 已提交
1739
    int32_t         numOfVg = taosArrayGetSize(pTopic->vgs);
1740 1741

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

1749
      int32_t vgStatus = atomic_val_compare_exchange_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE, TMQ_VG_STATUS__WAIT);
1750
      if (vgStatus == TMQ_VG_STATUS__WAIT) {
L
Liu Jicong 已提交
1751
        int32_t vgSkipCnt = atomic_add_fetch_32(&pVg->vgSkipCnt, 1);
1752
        tscTrace("consumer:0x%" PRIx64 " epoch %d wait poll-rsp, skip vgId:%d skip cnt %d", tmq->consumerId, tmq->epoch,
X
Xiaoyu Wang 已提交
1753
                 pVg->vgId, vgSkipCnt);
X
Xiaoyu Wang 已提交
1754
        continue;
L
temp  
Liu Jicong 已提交
1755 1756 1757 1758
#if 0
        if (skipCnt < 30000) {
          continue;
        } else {
1759
        tscDebug("consumer:0x%" PRIx64 ",skip vgId:%d skip too much reset", tmq->consumerId, pVg->vgId);
L
temp  
Liu Jicong 已提交
1760 1761
        }
#endif
X
Xiaoyu Wang 已提交
1762
      }
1763

L
Liu Jicong 已提交
1764
      atomic_store_32(&pVg->vgSkipCnt, 0);
1765 1766 1767
      int32_t code = doTmqPollImpl(tmq, pTopic, pVg, timeout);
      if (code != TSDB_CODE_SUCCESS) {
        return code;
D
dapan1121 已提交
1768
      }
X
Xiaoyu Wang 已提交
1769 1770
    }
  }
1771

1772
  tscDebug("consumer:0x%" PRIx64 " end to poll data", tmq->consumerId);
X
Xiaoyu Wang 已提交
1773 1774 1775
  return 0;
}

H
Haojun Liao 已提交
1776
static int32_t tmqHandleNoPollRsp(tmq_t* tmq, SMqRspWrapper* rspWrapper, bool* pReset) {
L
Liu Jicong 已提交
1777
  if (rspWrapper->tmqRspType == TMQ_MSG_TYPE__EP_RSP) {
L
fix  
Liu Jicong 已提交
1778
    /*printf("ep %d %d\n", rspMsg->head.epoch, tmq->epoch);*/
L
Liu Jicong 已提交
1779 1780
    if (rspWrapper->epoch > atomic_load_32(&tmq->epoch)) {
      SMqAskEpRspWrapper* pEpRspWrapper = (SMqAskEpRspWrapper*)rspWrapper;
L
Liu Jicong 已提交
1781
      SMqAskEpRsp*        rspMsg = &pEpRspWrapper->msg;
1782
      doUpdateLocalEp(tmq, rspWrapper->epoch, rspMsg);
L
temp  
Liu Jicong 已提交
1783
      /*tmqClearUnhandleMsg(tmq);*/
L
Liu Jicong 已提交
1784
      tDeleteSMqAskEpRsp(rspMsg);
X
Xiaoyu Wang 已提交
1785 1786
      *pReset = true;
    } else {
L
Liu Jicong 已提交
1787
      tmqFreeRspWrapper(rspWrapper);
X
Xiaoyu Wang 已提交
1788 1789 1790 1791 1792 1793 1794 1795
      *pReset = false;
    }
  } else {
    return -1;
  }
  return 0;
}

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

X
Xiaoyu Wang 已提交
1799
  while (1) {
1800 1801
    SMqRspWrapper* pRspWrapper = NULL;
    taosGetQitem(tmq->qall, (void**)&pRspWrapper);
1802

1803
    if (pRspWrapper == NULL) {
L
Liu Jicong 已提交
1804
      taosReadAllQitems(tmq->mqueue, tmq->qall);
1805 1806
      taosGetQitem(tmq->qall, (void**)&pRspWrapper);
      if (pRspWrapper == NULL) {
L
Liu Jicong 已提交
1807 1808
        return NULL;
      }
X
Xiaoyu Wang 已提交
1809 1810
    }

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

1813 1814
    if (pRspWrapper->tmqRspType == TMQ_MSG_TYPE__END_RSP) {
      taosFreeQitem(pRspWrapper);
L
Liu Jicong 已提交
1815
      terrno = TSDB_CODE_TQ_NO_COMMITTED_OFFSET;
H
Haojun Liao 已提交
1816
      tscError("consumer:0x%" PRIx64 " unexpected rsp from poll, code:%s", tmq->consumerId, tstrerror(terrno));
L
Liu Jicong 已提交
1817
      return NULL;
1818 1819
    } else if (pRspWrapper->tmqRspType == TMQ_MSG_TYPE__POLL_RSP) {
      SMqPollRspWrapper* pollRspWrapper = (SMqPollRspWrapper*)pRspWrapper;
H
Haojun Liao 已提交
1820

X
Xiaoyu Wang 已提交
1821
      int32_t     consumerEpoch = atomic_load_32(&tmq->epoch);
1822 1823 1824
      SMqDataRsp* pDataRsp = &pollRspWrapper->dataRsp;

      if (pDataRsp->head.epoch == consumerEpoch) {
wmmhello's avatar
wmmhello 已提交
1825 1826 1827 1828 1829 1830 1831 1832
        SMqClientVg* pVg = getVgInfo(tmq, pollRspWrapper->topicName, pollRspWrapper->vgId);
        pollRspWrapper->vgHandle = pVg;
        pollRspWrapper->topicHandle = getTopicInfo(tmq, pollRspWrapper->topicName);
        if(pollRspWrapper->vgHandle == NULL || pollRspWrapper->topicHandle == NULL){
          tscError("consumer:0x%" PRIx64 " get vg or topic error, topic:%s vgId:%d", tmq->consumerId,
                   pollRspWrapper->topicName, pollRspWrapper->vgId);
          return NULL;
        }
1833 1834 1835 1836
        // update the epset
        if (pollRspWrapper->pEpset != NULL) {
          SEp* pEp = GET_ACTIVE_EP(pollRspWrapper->pEpset);
          SEp* pOld = GET_ACTIVE_EP(&(pVg->epSet));
X
Xiaoyu Wang 已提交
1837 1838
          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);
1839 1840 1841
          pVg->epSet = *pollRspWrapper->pEpset;
        }

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

1847 1848 1849
        char buf[80];
        tFormatOffset(buf, 80, &pDataRsp->rspOffset);
        if (pDataRsp->blockNum == 0) {
X
Xiaoyu Wang 已提交
1850 1851 1852
          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);
1853
          pRspWrapper = tmqFreeRspWrapper(pRspWrapper);
1854
          pVg->emptyBlockReceiveTs = taosGetTimestampMs();
L
Liu Jicong 已提交
1855
          taosFreeQitem(pollRspWrapper);
1856
        } else {  // build rsp
X
Xiaoyu Wang 已提交
1857
          int64_t    numOfRows = 0;
1858
          SMqRspObj* pRsp = tmqBuildRspFromWrapper(pollRspWrapper, pVg, &numOfRows);
H
Haojun Liao 已提交
1859
          tmq->totalRows += numOfRows;
1860
          pVg->emptyBlockReceiveTs = 0;
H
Haojun Liao 已提交
1861
          tscDebug("consumer:0x%" PRIx64 " process poll rsp, vgId:%d, offset:%s, blocks:%d, rows:%" PRId64
H
Haojun Liao 已提交
1862
                   " vg total:%" PRId64 " total:%" PRId64 ", reqId:0x%" PRIx64,
H
Haojun Liao 已提交
1863
                   tmq->consumerId, pVg->vgId, buf, pDataRsp->blockNum, numOfRows, pVg->numOfRows, tmq->totalRows,
H
Haojun Liao 已提交
1864
                   pollRspWrapper->reqId);
1865 1866 1867
          taosFreeQitem(pollRspWrapper);
          return pRsp;
        }
X
Xiaoyu Wang 已提交
1868
      } else {
H
Haojun Liao 已提交
1869
        tscDebug("consumer:0x%" PRIx64 " vgId:%d msg discard since epoch mismatch: msg epoch %d, consumer epoch %d",
1870
                 tmq->consumerId, pollRspWrapper->vgId, pDataRsp->head.epoch, consumerEpoch);
1871
        pRspWrapper = tmqFreeRspWrapper(pRspWrapper);
L
Liu Jicong 已提交
1872 1873
        taosFreeQitem(pollRspWrapper);
      }
1874 1875
    } else if (pRspWrapper->tmqRspType == TMQ_MSG_TYPE__POLL_META_RSP) {
      SMqPollRspWrapper* pollRspWrapper = (SMqPollRspWrapper*)pRspWrapper;
L
Liu Jicong 已提交
1876
      int32_t            consumerEpoch = atomic_load_32(&tmq->epoch);
1877 1878 1879

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

L
Liu Jicong 已提交
1880
      if (pollRspWrapper->metaRsp.head.epoch == consumerEpoch) {
wmmhello's avatar
wmmhello 已提交
1881 1882 1883 1884 1885 1886 1887 1888
        SMqClientVg* pVg = getVgInfo(tmq, pollRspWrapper->topicName, pollRspWrapper->vgId);
        pollRspWrapper->vgHandle = pVg;
        pollRspWrapper->topicHandle = getTopicInfo(tmq, pollRspWrapper->topicName);
        if(pollRspWrapper->vgHandle == NULL || pollRspWrapper->topicHandle == NULL){
          tscError("consumer:0x%" PRIx64 " get vg or topic error, topic:%s vgId:%d", tmq->consumerId,
                   pollRspWrapper->topicName, pollRspWrapper->vgId);
          return NULL;
        }
1889 1890 1891
        if(pollRspWrapper->metaRsp.rspOffset.type != 0){    // if offset is validate
          pVg->currentOffset = pollRspWrapper->metaRsp.rspOffset;
        }
L
Liu Jicong 已提交
1892 1893
        atomic_store_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE);
        // build rsp
L
Liu Jicong 已提交
1894
        SMqMetaRspObj* pRsp = tmqBuildMetaRspFromWrapper(pollRspWrapper);
L
Liu Jicong 已提交
1895 1896 1897
        taosFreeQitem(pollRspWrapper);
        return pRsp;
      } else {
H
Haojun Liao 已提交
1898
        tscDebug("consumer:0x%" PRIx64 " vgId:%d msg discard since epoch mismatch: msg epoch %d, consumer epoch %d",
1899
                 tmq->consumerId, pollRspWrapper->vgId, pollRspWrapper->metaRsp.head.epoch, consumerEpoch);
1900
        pRspWrapper = tmqFreeRspWrapper(pRspWrapper);
L
Liu Jicong 已提交
1901
        taosFreeQitem(pollRspWrapper);
X
Xiaoyu Wang 已提交
1902
      }
1903 1904
    } else if (pRspWrapper->tmqRspType == TMQ_MSG_TYPE__TAOSX_RSP) {
      SMqPollRspWrapper* pollRspWrapper = (SMqPollRspWrapper*)pRspWrapper;
X
Xiaoyu Wang 已提交
1905
      int32_t            consumerEpoch = atomic_load_32(&tmq->epoch);
H
Haojun Liao 已提交
1906

L
Liu Jicong 已提交
1907
      if (pollRspWrapper->taosxRsp.head.epoch == consumerEpoch) {
wmmhello's avatar
wmmhello 已提交
1908 1909 1910 1911 1912 1913 1914 1915
        SMqClientVg* pVg = getVgInfo(tmq, pollRspWrapper->topicName, pollRspWrapper->vgId);
        pollRspWrapper->vgHandle = pVg;
        pollRspWrapper->topicHandle = getTopicInfo(tmq, pollRspWrapper->topicName);
        if(pollRspWrapper->vgHandle == NULL || pollRspWrapper->topicHandle == NULL){
          tscError("consumer:0x%" PRIx64 " get vg or topic error, topic:%s vgId:%d", tmq->consumerId,
                   pollRspWrapper->topicName, pollRspWrapper->vgId);
          return NULL;
        }
1916 1917 1918
        if(pollRspWrapper->taosxRsp.rspOffset.type != 0){    // if offset is validate
          pVg->currentOffset = pollRspWrapper->taosxRsp.rspOffset;
        }
L
Liu Jicong 已提交
1919
        atomic_store_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE);
H
Haojun Liao 已提交
1920

L
Liu Jicong 已提交
1921
        if (pollRspWrapper->taosxRsp.blockNum == 0) {
H
Haojun Liao 已提交
1922 1923
          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 已提交
1924
          pVg->emptyBlockReceiveTs = taosGetTimestampMs();
1925
          pRspWrapper = tmqFreeRspWrapper(pRspWrapper);
H
Haojun Liao 已提交
1926
          taosFreeQitem(pollRspWrapper);
L
Liu Jicong 已提交
1927
          continue;
H
Haojun Liao 已提交
1928
        } else {
X
Xiaoyu Wang 已提交
1929
          pVg->emptyBlockReceiveTs = 0;  // reset the ts
L
Liu Jicong 已提交
1930
        }
wmmhello's avatar
wmmhello 已提交
1931

L
Liu Jicong 已提交
1932
        // build rsp
X
Xiaoyu Wang 已提交
1933
        void*   pRsp = NULL;
1934
        int64_t numOfRows = 0;
L
Liu Jicong 已提交
1935
        if (pollRspWrapper->taosxRsp.createTableNum == 0) {
1936
          pRsp = tmqBuildRspFromWrapper(pollRspWrapper, pVg, &numOfRows);
L
Liu Jicong 已提交
1937
        } else {
wmmhello's avatar
wmmhello 已提交
1938 1939
          pRsp = tmqBuildTaosxRspFromWrapper(pollRspWrapper);
        }
H
Haojun Liao 已提交
1940

1941 1942
        tmq->totalRows += numOfRows;

H
Haojun Liao 已提交
1943 1944
        char buf[80];
        tFormatOffset(buf, 80, &pVg->currentOffset);
H
Haojun Liao 已提交
1945
        tscDebug("consumer:0x%" PRIx64 " process taosx poll rsp, vgId:%d, offset:%s, blocks:%d, rows:%" PRId64
X
Xiaoyu Wang 已提交
1946
                 ", vg total:%" PRId64 " total:%" PRId64 " reqId:0x%" PRIx64,
H
Haojun Liao 已提交
1947
                 tmq->consumerId, pVg->vgId, buf, pollRspWrapper->dataRsp.blockNum, numOfRows, pVg->numOfRows,
H
Haojun Liao 已提交
1948
                 tmq->totalRows, pollRspWrapper->reqId);
H
Haojun Liao 已提交
1949 1950

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

L
Liu Jicong 已提交
1953
      } else {
H
Haojun Liao 已提交
1954
        tscDebug("consumer:0x%" PRIx64 " vgId:%d msg discard since epoch mismatch: msg epoch %d, consumer epoch %d",
1955
                 tmq->consumerId, pollRspWrapper->vgId, pollRspWrapper->taosxRsp.head.epoch, consumerEpoch);
1956
        pRspWrapper = tmqFreeRspWrapper(pRspWrapper);
L
Liu Jicong 已提交
1957 1958
        taosFreeQitem(pollRspWrapper);
      }
X
Xiaoyu Wang 已提交
1959
    } else {
H
Haojun Liao 已提交
1960 1961
      tscDebug("consumer:0x%" PRIx64 " not data msg received", tmq->consumerId);

X
Xiaoyu Wang 已提交
1962
      bool reset = false;
1963 1964
      tmqHandleNoPollRsp(tmq, pRspWrapper, &reset);
      taosFreeQitem(pRspWrapper);
X
Xiaoyu Wang 已提交
1965
      if (pollIfReset && reset) {
1966
        tscDebug("consumer:0x%" PRIx64 ", reset and repoll", tmq->consumerId);
1967
        tmqPollImpl(tmq, timeout);
X
Xiaoyu Wang 已提交
1968 1969 1970 1971 1972
      }
    }
  }
}

1973
TAOS_RES* tmq_consumer_poll(tmq_t* tmq, int64_t timeout) {
L
Liu Jicong 已提交
1974 1975
  void*   rspObj;
  int64_t startTime = taosGetTimestampMs();
L
Liu Jicong 已提交
1976

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

1980 1981 1982
#if 0
  tmqHandleAllDelayedTask(tmq);
  tmqPollImpl(tmq, timeout);
1983
  rspObj = tmqHandleAllRsp(tmq, timeout, false);
L
Liu Jicong 已提交
1984 1985
  if (rspObj) {
    return (TAOS_RES*)rspObj;
L
fix  
Liu Jicong 已提交
1986
  }
1987
#endif
X
Xiaoyu Wang 已提交
1988

1989
  // in no topic status, delayed task also need to be processed
L
Liu Jicong 已提交
1990
  if (atomic_load_8(&tmq->status) == TMQ_CONSUMER_STATUS__INIT) {
1991
    tscDebug("consumer:0x%" PRIx64 " poll return since consumer is init", tmq->consumerId);
1992
    taosMsleep(500);  //     sleep for a while
1993 1994 1995
    return NULL;
  }

wmmhello's avatar
wmmhello 已提交
1996
  while (atomic_load_8(&tmq->status) == TMQ_CONSUMER_STATUS__RECOVER) {
L
Liu Jicong 已提交
1997
    int32_t retryCnt = 0;
1998
    while (TSDB_CODE_MND_CONSUMER_NOT_READY == doAskEp(tmq)) {
H
Haojun Liao 已提交
1999
      if (retryCnt++ > 40) {
L
Liu Jicong 已提交
2000 2001
        return NULL;
      }
2002

H
Haojun Liao 已提交
2003
      tscDebug("consumer:0x%" PRIx64 " not ready, retry:%d/40 in 500ms", tmq->consumerId, retryCnt);
L
Liu Jicong 已提交
2004 2005 2006 2007
      taosMsleep(500);
    }
  }

X
Xiaoyu Wang 已提交
2008
  while (1) {
L
Liu Jicong 已提交
2009
    tmqHandleAllDelayedTask(tmq);
2010

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

2015
    rspObj = tmqHandleAllRsp(tmq, timeout, false);
L
Liu Jicong 已提交
2016
    if (rspObj) {
2017
      tscDebug("consumer:0x%" PRIx64 " return rsp %p", tmq->consumerId, rspObj);
L
Liu Jicong 已提交
2018
      return (TAOS_RES*)rspObj;
L
Liu Jicong 已提交
2019
    } else if (terrno == TSDB_CODE_TQ_NO_COMMITTED_OFFSET) {
2020
      tscDebug("consumer:0x%" PRIx64 " return null since no committed offset", tmq->consumerId);
L
Liu Jicong 已提交
2021
      return NULL;
X
Xiaoyu Wang 已提交
2022
    }
2023

2024
    if (timeout >= 0) {
L
Liu Jicong 已提交
2025
      int64_t currentTime = taosGetTimestampMs();
2026 2027 2028
      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 已提交
2029
                 tmq->consumerId, tmq->epoch, startTime, currentTime);
X
Xiaoyu Wang 已提交
2030 2031
        return NULL;
      }
2032
      tsem_timewait(&tmq->rspSem, (timeout - elapsedTime));
L
Liu Jicong 已提交
2033 2034
    } else {
      // use tsem_timewait instead of tsem_wait to avoid unexpected stuck
L
Liu Jicong 已提交
2035
      tsem_timewait(&tmq->rspSem, 1000);
X
Xiaoyu Wang 已提交
2036 2037 2038 2039
    }
  }
}

2040 2041 2042 2043 2044 2045 2046 2047 2048 2049 2050 2051 2052 2053
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);
2054
    }
2055
  }
2056

2057 2058
  tscDebug("consumer:0x%" PRIx64 " rows dist end", pTmq->consumerId);
}
2059

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

2064 2065 2066 2067 2068 2069
  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;
2070 2071 2072
      }
    }

L
Liu Jicong 已提交
2073
    int32_t     retryCnt = 0;
2074
    tmq_list_t* lst = tmq_list_new();
L
Liu Jicong 已提交
2075
    while (1) {
2076
      int32_t rsp = tmq_subscribe(tmq, lst);
L
Liu Jicong 已提交
2077 2078 2079 2080 2081 2082 2083 2084
      if (rsp != TSDB_CODE_MND_CONSUMER_NOT_READY || retryCnt > 5) {
        break;
      } else {
        retryCnt++;
        taosMsleep(500);
      }
    }

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

2090
  taosRemoveRef(tmqMgmt.rsetId, tmq->refId);
L
Liu Jicong 已提交
2091
  return 0;
2092
}
L
Liu Jicong 已提交
2093

L
Liu Jicong 已提交
2094 2095
const char* tmq_err2str(int32_t err) {
  if (err == 0) {
L
Liu Jicong 已提交
2096
    return "success";
L
Liu Jicong 已提交
2097
  } else if (err == -1) {
L
Liu Jicong 已提交
2098 2099 2100
    return "fail";
  } else {
    return tstrerror(err);
L
Liu Jicong 已提交
2101 2102
  }
}
L
Liu Jicong 已提交
2103

L
Liu Jicong 已提交
2104 2105 2106 2107 2108
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;
2109 2110
  } else if (TD_RES_TMQ_METADATA(res)) {
    return TMQ_RES_METADATA;
L
Liu Jicong 已提交
2111 2112 2113 2114 2115
  } else {
    return TMQ_RES_INVALID;
  }
}

L
Liu Jicong 已提交
2116
const char* tmq_get_topic_name(TAOS_RES* res) {
L
Liu Jicong 已提交
2117 2118
  if (TD_RES_TMQ(res)) {
    SMqRspObj* pRspObj = (SMqRspObj*)res;
L
Liu Jicong 已提交
2119
    return strchr(pRspObj->topic, '.') + 1;
L
Liu Jicong 已提交
2120 2121 2122
  } else if (TD_RES_TMQ_META(res)) {
    SMqMetaRspObj* pMetaRspObj = (SMqMetaRspObj*)res;
    return strchr(pMetaRspObj->topic, '.') + 1;
2123 2124 2125
  } else if (TD_RES_TMQ_METADATA(res)) {
    SMqTaosxRspObj* pRspObj = (SMqTaosxRspObj*)res;
    return strchr(pRspObj->topic, '.') + 1;
L
Liu Jicong 已提交
2126 2127 2128 2129 2130
  } else {
    return NULL;
  }
}

L
Liu Jicong 已提交
2131 2132 2133 2134
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 已提交
2135 2136 2137
  } else if (TD_RES_TMQ_META(res)) {
    SMqMetaRspObj* pMetaRspObj = (SMqMetaRspObj*)res;
    return strchr(pMetaRspObj->db, '.') + 1;
2138 2139 2140
  } else if (TD_RES_TMQ_METADATA(res)) {
    SMqTaosxRspObj* pRspObj = (SMqTaosxRspObj*)res;
    return strchr(pRspObj->db, '.') + 1;
L
Liu Jicong 已提交
2141 2142 2143 2144 2145
  } else {
    return NULL;
  }
}

L
Liu Jicong 已提交
2146 2147 2148 2149
int32_t tmq_get_vgroup_id(TAOS_RES* res) {
  if (TD_RES_TMQ(res)) {
    SMqRspObj* pRspObj = (SMqRspObj*)res;
    return pRspObj->vgId;
L
Liu Jicong 已提交
2150 2151 2152
  } else if (TD_RES_TMQ_META(res)) {
    SMqMetaRspObj* pMetaRspObj = (SMqMetaRspObj*)res;
    return pMetaRspObj->vgId;
2153
  } else if (TD_RES_TMQ_METADATA(res)) {
2154 2155
    SMqTaosxRspObj* pRspObj = (SMqTaosxRspObj*)res;
    return pRspObj->vgId;
L
Liu Jicong 已提交
2156 2157 2158 2159
  } else {
    return -1;
  }
}
L
Liu Jicong 已提交
2160 2161 2162 2163 2164 2165 2166 2167

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;
    }
2168
    return (const char*)taosArrayGetP(pRspObj->rsp.blockTbName, pRspObj->resIter);
2169 2170
  } else if (TD_RES_TMQ_METADATA(res)) {
    SMqTaosxRspObj* pRspObj = (SMqTaosxRspObj*)res;
L
Liu Jicong 已提交
2171 2172 2173
    if (!pRspObj->rsp.withTbName || pRspObj->rsp.blockTbName == NULL || pRspObj->resIter < 0 ||
        pRspObj->resIter >= pRspObj->rsp.blockNum) {
      return NULL;
2174
    }
L
Liu Jicong 已提交
2175 2176
    return (const char*)taosArrayGetP(pRspObj->rsp.blockTbName, pRspObj->resIter);
  }
L
Liu Jicong 已提交
2177 2178
  return NULL;
}
2179

2180 2181 2182 2183 2184 2185
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 已提交
2186 2187
}

2188 2189 2190 2191
static void commitCallBackFn(tmq_t *pTmq, int32_t code, void* param) {
  SSyncCommitInfo* pInfo = (SSyncCommitInfo*) param;
  pInfo->code = code;
  tsem_post(&pInfo->sem);
2192
}
2193

2194 2195 2196 2197 2198 2199 2200 2201 2202
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 已提交
2203 2204
  } else {
    asyncCommitOffset(tmq, pRes, commitCallBackFn, pInfo);
2205 2206
  }

2207 2208
  tsem_wait(&pInfo->sem);
  code = pInfo->code;
H
Haojun Liao 已提交
2209 2210

  tsem_destroy(&pInfo->sem);
2211 2212
  taosMemoryFree(pInfo);

X
Xiaoyu Wang 已提交
2213
  tscDebug("consumer:0x%" PRIx64 " sync commit done, code:%s", tmq->consumerId, tstrerror(code));
2214 2215 2216 2217 2218 2219 2220 2221 2222 2223 2224 2225
  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);
2226
    doUpdateLocalEp(pTmq, head->epoch, &rsp);
2227 2228 2229
    tDeleteSMqAskEpRsp(&rsp);
  }

H
Haojun Liao 已提交
2230
  tsem_post(&pInfo->sem);
2231 2232 2233 2234 2235 2236 2237 2238 2239 2240 2241 2242 2243 2244 2245 2246 2247 2248 2249 2250 2251 2252 2253 2254 2255 2256
}

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 已提交
2257
  tsem_init(&pInfo->sem, 0, 0);
2258 2259

  asyncAskEp(pTmq, updateEpCallbackFn, pInfo);
H
Haojun Liao 已提交
2260
  tsem_wait(&pInfo->sem);
2261 2262

  int32_t code = pInfo->code;
H
Haojun Liao 已提交
2263
  tsem_destroy(&pInfo->sem);
2264 2265 2266 2267 2268
  taosMemoryFree(pInfo);
  return code;
}

void asyncAskEp(tmq_t* pTmq, __tmq_askep_fn_t askEpFn, void* param) {
2269
  SMqAskEpReq req = {0};
2270 2271 2272
  req.consumerId = pTmq->consumerId;
  req.epoch = pTmq->epoch;
  strcpy(req.cgroup, pTmq->groupId);
2273 2274 2275

  int32_t tlen = tSerializeSMqAskEpReq(NULL, 0, &req);
  if (tlen < 0) {
2276 2277 2278
    tscError("consumer:0x%" PRIx64 ", tSerializeSMqAskEpReq failed", pTmq->consumerId);
    askEpFn(pTmq, TSDB_CODE_INVALID_PARA, NULL, param);
    return;
2279 2280 2281 2282
  }

  void* pReq = taosMemoryCalloc(1, tlen);
  if (pReq == NULL) {
2283 2284 2285
    tscError("consumer:0x%" PRIx64 ", failed to malloc askEpReq msg, size:%d", pTmq->consumerId, tlen);
    askEpFn(pTmq, TSDB_CODE_OUT_OF_MEMORY, NULL, param);
    return;
2286 2287 2288
  }

  if (tSerializeSMqAskEpReq(pReq, tlen, &req) < 0) {
2289
    tscError("consumer:0x%" PRIx64 ", tSerializeSMqAskEpReq %d failed", pTmq->consumerId, tlen);
2290
    taosMemoryFree(pReq);
2291 2292 2293

    askEpFn(pTmq, TSDB_CODE_INVALID_PARA, NULL, param);
    return;
2294 2295 2296 2297
  }

  SMqAskEpCbParam* pParam = taosMemoryCalloc(1, sizeof(SMqAskEpCbParam));
  if (pParam == NULL) {
2298
    tscError("consumer:0x%" PRIx64 ", failed to malloc subscribe param", pTmq->consumerId);
2299
    taosMemoryFree(pReq);
2300 2301 2302

    askEpFn(pTmq, TSDB_CODE_OUT_OF_MEMORY, NULL, param);
    return;
2303 2304
  }

2305 2306 2307 2308
  pParam->refId = pTmq->refId;
  pParam->epoch = pTmq->epoch;
  pParam->pUserFn = askEpFn;
  pParam->pParam = param;
2309 2310 2311 2312 2313

  SMsgSendInfo* sendInfo = taosMemoryCalloc(1, sizeof(SMsgSendInfo));
  if (sendInfo == NULL) {
    taosMemoryFree(pParam);
    taosMemoryFree(pReq);
2314 2315
    askEpFn(pTmq, TSDB_CODE_OUT_OF_MEMORY, NULL, param);
    return;
2316 2317
  }

X
Xiaoyu Wang 已提交
2318
  sendInfo->msgInfo = (SDataBuf){.pData = pReq, .len = tlen, .handle = NULL};
2319 2320 2321 2322

  sendInfo->requestId = generateRequestId();
  sendInfo->requestObjRefId = 0;
  sendInfo->param = pParam;
2323
  sendInfo->fp = askEpCallbackFn;
2324 2325
  sendInfo->msgType = TDMT_MND_TMQ_ASK_EP;

2326 2327
  SEpSet epSet = getEpSet_s(&pTmq->pTscObj->pAppInfo->mgmtEp);
  tscDebug("consumer:0x%" PRIx64 " ask ep from mnode, reqId:0x%" PRIx64, pTmq->consumerId, sendInfo->requestId);
2328 2329

  int64_t transporterId = 0;
2330
  asyncSendMsgToServer(pTmq->pTscObj->pAppInfo->pTransporter, &epSet, &transporterId, sendInfo);
2331 2332 2333 2334 2335 2336 2337
}

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

int32_t tmqCommitDone(SMqCommitCbParamSet* pParamSet) {
2338 2339 2340
  int64_t refId = pParamSet->refId;

  tmq_t* tmq = taosAcquireRef(tmqMgmt.rsetId, refId);
2341 2342 2343 2344 2345 2346 2347
  if (tmq == NULL) {
    taosMemoryFree(pParamSet);
    terrno = TSDB_CODE_TMQ_CONSUMER_CLOSED;
    return -1;
  }

  // if no more waiting rsp
2348
  pParamSet->callbackFn(tmq, pParamSet->code, pParamSet->userParam);
2349
  taosMemoryFree(pParamSet);
2350 2351

  taosReleaseRef(tmqMgmt.rsetId, refId);
2352
  return 0;
2353 2354
}

2355
void commitRspCountDown(SMqCommitCbParamSet* pParamSet, int64_t consumerId, const char* pTopic, int32_t vgId) {
2356 2357
  int32_t waitingRspNum = atomic_sub_fetch_32(&pParamSet->waitingRspNum, 1);
  if (waitingRspNum == 0) {
H
Haojun Liao 已提交
2358 2359
    tscDebug("consumer:0x%" PRIx64 " topic:%s vgId:%d all commit-rsp received, commit completed", consumerId, pTopic,
             vgId);
2360
    tmqCommitDone(pParamSet);
H
Haojun Liao 已提交
2361 2362 2363
  } else {
    tscDebug("consumer:0x%" PRIx64 " topic:%s vgId:%d commit-rsp received, remain:%d", consumerId, pTopic, vgId,
             waitingRspNum);
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

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