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

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

27
#define EMPTY_BLOCK_POLL_IDLE_DURATION  10
28
#define DEFAULT_AUTO_COMMIT_INTERVAL    5000
29

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

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

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

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

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

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

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

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

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

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

114 115 116 117
typedef struct SAskEpInfo {
  int32_t code;
} SAskEpInfo;

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

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

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

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

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

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

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

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

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

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

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

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

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

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

236 237 238
  return conf;
}

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

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

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

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

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

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

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

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

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

326 327
  if (strcasecmp(key, "enable.heartbeat.background") == 0) {
    if (strcasecmp(value, "true") == 0) {
328
      conf->hbBgEnable = true;
L
Liu Jicong 已提交
329
      return TMQ_CONF_OK;
330
    } else if (strcasecmp(value, "false") == 0) {
331
      conf->hbBgEnable = false;
L
Liu Jicong 已提交
332 333 334 335 336 337
      return TMQ_CONF_OK;
    } else {
      return TMQ_CONF_INVALID;
    }
  }

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

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

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

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

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

L
Liu Jicong 已提交
362
  return TMQ_CONF_UNKNOWN;
363 364 365
}

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

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

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

L
Liu Jicong 已提交
385 386 387 388 389 390 391 392 393 394
int32_t tmq_list_get_size(const tmq_list_t* list) {
  const SArray* container = &list->container;
  return taosArrayGetSize(container);
}

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

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

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

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

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

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

426 427 428 429 430 431 432 433 434 435 436 437 438 439 440 441 442 443 444 445 446 447 448 449
//  if (code != TSDB_CODE_SUCCESS) { // if commit offset failed, let's try again
//    taosThreadMutexLock(&pParam->pTmq->lock);
//    int32_t numOfVgroups, index;
//    SMqClientVg* pVg = foundClientVg(pParam->pTmq->clientTopics, pParam->topicName, pParam->vgId, &index, &numOfVgroups);
//    if (pVg == NULL) {
//      tscDebug("consumer:0x%" PRIx64
//               " subKey:%s vgId:%d commit failed, code:%s has been transferred to other consumer, no need retry ordinal:%d/%d",
//               pParam->pTmq->consumerId, pParam->pOffset->subKey, pParam->vgId, tstrerror(code), index + 1, numOfVgroups);
//    } else { // let's retry the commit
//      int32_t code1 = doSendCommitMsg(pParam->pTmq, pVg, pParam->topicName, pParamSet, index, numOfVgroups);
//      if (code1 != TSDB_CODE_SUCCESS) {  // retry failed.
//        tscError("consumer:0x%" PRIx64 " topic:%s vgId:%d offset:%" PRId64
//                 " retry failed, ignore this commit. code:%s ordinal:%d/%d",
//                 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);
//
450
//    commitRspCountDown(pParamSet, pParam->pTmq->consumerId, pParam->topicName, pParam->vgId);
451 452 453 454
//    return 0;
//  }
//
//  // todo replace the pTmq with refId
455

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

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

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

L
Liu Jicong 已提交
471
  pOffset->val = pVg->currentOffset;
472

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

633 634 635
  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 已提交
636

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

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

  pParamSet->refId = tmq->refId;
  pParamSet->epoch = tmq->epoch;
657
  pParamSet->userCb = pCommitFp;
658 659
  pParamSet->userParam = userParam;

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

  return NULL;
L
Liu Jicong 已提交
885 886
}

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

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

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

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

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

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

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

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

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

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

  taosFreeQall(tmq->qall);
971
  tsem_destroy(&tmq->rspSem);
H
Haojun Liao 已提交
972
  taosThreadMutexDestroy(&tmq->lock);
L
Liu Jicong 已提交
973

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

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

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

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

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

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

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

L
Liu Jicong 已提交
1026 1027
  // init status
  pTmq->status = TMQ_CONSUMER_STATUS__INIT;
L
Liu Jicong 已提交
1028 1029
  pTmq->pollCnt = 0;
  pTmq->epoch = 0;
L
Liu Jicong 已提交
1030 1031
  /*pTmq->epStatus = 0;*/
  /*pTmq->epSkipCnt = 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 1043
  pTmq->resetOffsetCfg = conf->resetOffset;

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

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

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

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

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

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

1075 1076 1077
  char         buf[80] = {0};
  STqOffsetVal offset = {.type = pTmq->resetOffsetCfg};
  tFormatOffset(buf, tListLen(buf), &offset);
1078 1079
  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,
1080
          pTmq->hbBgEnable);
L
Liu Jicong 已提交
1081

1082
  return pTmq;
1083

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

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

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

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

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

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

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

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

    taosArrayPush(req.topicNames, &topicFName);
1124 1125
  }

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

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

1134 1135 1136
  void* abuf = buf;
  tSerializeSCMSubscribeReq(&abuf, &req);

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

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

1149 1150 1151
  if (tsem_init(&param.rspSem, 0, 0) != 0) {
    goto FAIL;
  }
L
Liu Jicong 已提交
1152 1153

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

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

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

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

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

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

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

L
Liu Jicong 已提交
1182
  int32_t retryCnt = 0;
1183
  while (TSDB_CODE_MND_CONSUMER_NOT_READY == doAskEp(tmq)) {
1184
    if (retryCnt++ > MAX_RETRY_COUNT) {
L
Liu Jicong 已提交
1185 1186
      goto FAIL;
    }
1187

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

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

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

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

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

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

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

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

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

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

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

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

L
Liu Jicong 已提交
1246
    if (pMsg->pData) taosMemoryFree(pMsg->pData);
dengyihao's avatar
dengyihao 已提交
1247 1248
    if (pMsg->pEpSet) taosMemoryFree(pMsg->pEpSet);

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

L
Liu Jicong 已提交
1262 1263 1264
      pRspWrapper->tmqRspType = TMQ_MSG_TYPE__END_RSP;
      taosWriteQitem(tmq->mqueue, pRspWrapper);
    }
H
Haojun Liao 已提交
1265

L
fix txn  
Liu Jicong 已提交
1266
    goto CREATE_MSG_FAIL;
L
Liu Jicong 已提交
1267 1268
  }

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

1276
    tsem_post(&tmq->rspSem);
1277 1278
    taosReleaseRef(tmqMgmt.rsetId, refId);

L
Liu Jicong 已提交
1279
    taosMemoryFree(pMsg->pData);
dengyihao's avatar
dengyihao 已提交
1280
    taosMemoryFree(pMsg->pEpSet);
X
Xiaoyu Wang 已提交
1281 1282 1283 1284
    return 0;
  }

  if (msgEpoch != tmqEpoch) {
H
Haojun Liao 已提交
1285 1286
    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 已提交
1287 1288
  }

L
Liu Jicong 已提交
1289 1290 1291
  // handle meta rsp
  int8_t rspType = ((SMqRspHead*)pMsg->pData)->mqMsgType;

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

L
Liu Jicong 已提交
1300
  pRspWrapper->tmqRspType = rspType;
L
Liu Jicong 已提交
1301 1302
  pRspWrapper->vgHandle = pVg;
  pRspWrapper->topicHandle = pTopic;
H
Haojun Liao 已提交
1303
  pRspWrapper->reqId = requestId;
1304
  pRspWrapper->pEpset = pMsg->pEpSet;
1305
  pRspWrapper->vgId = pVg->vgId;
L
Liu Jicong 已提交
1306

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

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

L
Liu Jicong 已提交
1335
  taosMemoryFree(pMsg->pData);
H
Haojun Liao 已提交
1336
  taosWriteQitem(tmq->mqueue, pRspWrapper);
L
Liu Jicong 已提交
1337

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

1341
  tsem_post(&tmq->rspSem);
1342 1343
  taosReleaseRef(tmqMgmt.rsetId, refId);

L
Liu Jicong 已提交
1344
  return 0;
H
Haojun Liao 已提交
1345

L
fix txn  
Liu Jicong 已提交
1346
CREATE_MSG_FAIL:
L
Liu Jicong 已提交
1347
  if (epoch == tmq->epoch) {
L
Liu Jicong 已提交
1348 1349
    atomic_store_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE);
  }
H
Haojun Liao 已提交
1350

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

L
Liu Jicong 已提交
1354
  return -1;
1355 1356
}

H
Haojun Liao 已提交
1357 1358 1359 1360 1361
typedef struct SVgroupSaveInfo {
  STqOffsetVal offset;
  int64_t      numOfRows;
} SVgroupSaveInfo;

H
Haojun Liao 已提交
1362 1363 1364 1365 1366 1367 1368 1369 1370 1371 1372 1373 1374 1375 1376 1377 1378
static void initClientTopicFromRsp(SMqClientTopic* pTopic, SMqSubTopicEp* pTopicEp, SHashObj* pVgOffsetHashMap,
                                   tmq_t* tmq) {
  pTopic->schema = pTopicEp->schema;
  pTopicEp->schema.nCols = 0;
  pTopicEp->schema.pSchema = NULL;

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

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

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

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

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

H
Haojun Liao 已提交
1383
    int64_t numOfRows = 0;
H
Haojun Liao 已提交
1384
    STqOffsetVal  offsetNew = {.type = tmq->resetOffsetCfg};
H
Haojun Liao 已提交
1385 1386 1387
    if (pInfo != NULL) {
      offsetNew = pInfo->offset;
      numOfRows = pInfo->numOfRows;
H
Haojun Liao 已提交
1388 1389 1390 1391 1392 1393 1394 1395 1396
    }

    SMqClientVg clientVg = {
        .pollCnt = 0,
        .currentOffset = offsetNew,
        .vgId = pVgEp->vgId,
        .epSet = pVgEp->epSet,
        .vgStatus = TMQ_VG_STATUS__IDLE,
        .vgSkipCnt = 0,
H
Haojun Liao 已提交
1397
        .emptyBlockReceiveTs = 0,
H
Haojun Liao 已提交
1398
        .numOfRows = numOfRows,
H
Haojun Liao 已提交
1399 1400 1401 1402 1403 1404 1405 1406 1407 1408 1409 1410 1411 1412 1413 1414
    };

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

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

  taosArrayDestroy(pTopic->vgs);
}

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

1417
  int32_t topicNumCur = taosArrayGetSize(tmq->clientTopics);
1418
  int32_t topicNumGet = taosArrayGetSize(pRsp->topics);
1419

X
Xiaoyu Wang 已提交
1420 1421
  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",
1422
           tmq->consumerId, tmq->epoch, epoch, topicNumGet, topicNumCur);
1423 1424 1425 1426 1427 1428

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

H
Haojun Liao 已提交
1429 1430
  SHashObj* pVgOffsetHashMap = taosHashInit(64, MurmurHash3_32, false, HASH_NO_LOCK);
  if (pVgOffsetHashMap == NULL) {
1431 1432 1433
    taosArrayDestroy(newTopics);
    return false;
  }
1434

H
Haojun Liao 已提交
1435
  // todo extract method
1436 1437 1438 1439 1440
  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);
1441
      tscDebug("consumer:0x%" PRIx64 ", new vg num: %d", tmq->consumerId, vgNumCur);
1442 1443
      for (int32_t j = 0; j < vgNumCur; j++) {
        SMqClientVg* pVgCur = taosArrayGet(pTopicCur->vgs, j);
H
Haojun Liao 已提交
1444 1445
        makeTopicVgroupKey(vgKey, pTopicCur->topicName, pVgCur->vgId);

L
Liu Jicong 已提交
1446
        char buf[80];
L
Liu Jicong 已提交
1447
        tFormatOffset(buf, 80, &pVgCur->currentOffset);
H
Haojun Liao 已提交
1448
        tscDebug("consumer:0x%" PRIx64 ", epoch:%d vgId:%d vgKey:%s, offset:%s", tmq->consumerId, epoch,
L
Liu Jicong 已提交
1449
                 pVgCur->vgId, vgKey, buf);
H
Haojun Liao 已提交
1450 1451 1452

        SVgroupSaveInfo info = {.offset = pVgCur->currentOffset, .numOfRows = pVgCur->numOfRows};
        taosHashPut(pVgOffsetHashMap, vgKey, strlen(vgKey), &info, sizeof(SVgroupSaveInfo));
1453 1454 1455 1456 1457 1458 1459
      }
    }
  }

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

H
Haojun Liao 已提交
1464 1465 1466
  taosHashCleanup(pVgOffsetHashMap);

  taosThreadMutexLock(&tmq->lock);
1467
  // destroy current buffered existed topics info
1468
  if (tmq->clientTopics) {
H
Haojun Liao 已提交
1469
    taosArrayDestroyEx(tmq->clientTopics, freeClientVgInfo);
X
Xiaoyu Wang 已提交
1470
  }
1471

H
Haojun Liao 已提交
1472 1473
  tmq->clientTopics = newTopics;
  taosThreadMutexUnlock(&tmq->lock);
1474

H
Haojun Liao 已提交
1475 1476
  int8_t flag = (topicNumGet == 0)? TMQ_CONSUMER_STATUS__NO_TOPIC:TMQ_CONSUMER_STATUS__READY;
  atomic_store_8(&tmq->status, flag);
X
Xiaoyu Wang 已提交
1477
  atomic_store_32(&tmq->epoch, epoch);
H
Haojun Liao 已提交
1478

1479
  tscDebug("consumer:0x%" PRIx64 " update topic info completed", tmq->consumerId);
X
Xiaoyu Wang 已提交
1480 1481 1482
  return set;
}

1483
int32_t askEpCallbackFn(void* param, SDataBuf* pMsg, int32_t code) {
1484
  SMqAskEpCbParam* pParam = (SMqAskEpCbParam*)param;
1485
  tmq_t* tmq = taosAcquireRef(tmqMgmt.rsetId, pParam->refId);
1486 1487

  if (tmq == NULL) {
1488 1489 1490
    terrno = TSDB_CODE_TMQ_CONSUMER_CLOSED;
    pParam->pUserFn(tmq, terrno, NULL, pParam->pParam);

1491
    taosMemoryFree(pMsg->pData);
dengyihao's avatar
dengyihao 已提交
1492
    taosMemoryFree(pMsg->pEpSet);
1493 1494
    taosMemoryFree(pParam);
    return terrno;
1495 1496
  }

H
Haojun Liao 已提交
1497
  if (code != TSDB_CODE_SUCCESS) {
1498 1499 1500 1501 1502 1503 1504 1505 1506
    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;
1507
  }
L
Liu Jicong 已提交
1508

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

1518 1519 1520 1521 1522 1523 1524 1525
    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 已提交
1526
  } else {
1527 1528 1529
    tscDebug("consumer:0x%" PRIx64 ", recv ep, msg epoch %d, current epoch %d, update local ep", tmq->consumerId,
             head->epoch, epoch);
    pParam->pUserFn(tmq, code, pMsg, pParam->pParam);
1530
  }
L
Liu Jicong 已提交
1531

1532 1533
  taosReleaseRef(tmqMgmt.rsetId, pParam->refId);

dengyihao's avatar
dengyihao 已提交
1534
  taosMemoryFree(pMsg->pEpSet);
L
Liu Jicong 已提交
1535
  taosMemoryFree(pMsg->pData);
1536
  taosMemoryFree(pParam);
L
Liu Jicong 已提交
1537
  return code;
1538 1539
}

L
Liu Jicong 已提交
1540
void tmqBuildConsumeReqImpl(SMqPollReq* pReq, tmq_t* tmq, int64_t timeout, SMqClientTopic* pTopic, SMqClientVg* pVg) {
L
Liu Jicong 已提交
1541 1542 1543 1544
  int32_t groupLen = strlen(tmq->groupId);
  memcpy(pReq->subKey, tmq->groupId, groupLen);
  pReq->subKey[groupLen] = TMQ_SEPARATOR;
  strcpy(pReq->subKey + groupLen + 1, pTopic->topicName);
1545

1546
  pReq->withTbName = tmq->withTbName;
L
Liu Jicong 已提交
1547
  pReq->consumerId = tmq->consumerId;
1548
  pReq->timeout = timeout;
X
Xiaoyu Wang 已提交
1549
  pReq->epoch = tmq->epoch;
L
Liu Jicong 已提交
1550
  /*pReq->currentOffset = reqOffset;*/
L
Liu Jicong 已提交
1551
  pReq->reqOffset = pVg->currentOffset;
D
dapan1121 已提交
1552
  pReq->head.vgId = pVg->vgId;
1553 1554
  pReq->useSnapshot = tmq->useSnapshot;
  pReq->reqId = generateRequestId();
1555 1556
}

L
Liu Jicong 已提交
1557 1558
SMqMetaRspObj* tmqBuildMetaRspFromWrapper(SMqPollRspWrapper* pWrapper) {
  SMqMetaRspObj* pRspObj = taosMemoryCalloc(1, sizeof(SMqMetaRspObj));
L
Liu Jicong 已提交
1559
  pRspObj->resType = RES_TYPE__TMQ_META;
L
Liu Jicong 已提交
1560 1561 1562 1563 1564 1565 1566 1567
  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;
}

1568
SMqRspObj* tmqBuildRspFromWrapper(SMqPollRspWrapper* pWrapper, SMqClientVg* pVg, int64_t* numOfRows) {
L
Liu Jicong 已提交
1569 1570
  SMqRspObj* pRspObj = taosMemoryCalloc(1, sizeof(SMqRspObj));
  pRspObj->resType = RES_TYPE__TMQ;
1571

1572
  (*numOfRows) = 0;
L
Liu Jicong 已提交
1573 1574
  tstrncpy(pRspObj->topic, pWrapper->topicHandle->topicName, TSDB_TOPIC_FNAME_LEN);
  tstrncpy(pRspObj->db, pWrapper->topicHandle->db, TSDB_DB_FNAME_LEN);
1575

L
Liu Jicong 已提交
1576
  pRspObj->vgId = pWrapper->vgHandle->vgId;
L
Liu Jicong 已提交
1577
  pRspObj->resIter = -1;
L
Liu Jicong 已提交
1578
  memcpy(&pRspObj->rsp, &pWrapper->dataRsp, sizeof(SMqDataRsp));
L
Liu Jicong 已提交
1579

L
Liu Jicong 已提交
1580 1581
  pRspObj->resInfo.totalRows = 0;
  pRspObj->resInfo.precision = TSDB_TIME_PRECISION_MILLI;
1582

L
Liu Jicong 已提交
1583
  if (!pWrapper->dataRsp.withSchema) {
L
Liu Jicong 已提交
1584 1585
    setResSchemaInfo(&pRspObj->resInfo, pWrapper->topicHandle->schema.pSchema, pWrapper->topicHandle->schema.nCols);
  }
L
Liu Jicong 已提交
1586

1587 1588 1589 1590 1591
  // extract the rows in this data packet
  for(int32_t i = 0; i < pRspObj->rsp.blockNum; ++i) {
    SRetrieveTableRsp* pRetrieve = (SRetrieveTableRsp*)taosArrayGetP(pRspObj->rsp.blockData, i);
    int64_t rows = htobe64(pRetrieve->numOfRows);
    pVg->numOfRows += rows;
1592
    (*numOfRows) += rows;
1593 1594
  }

L
Liu Jicong 已提交
1595
  return pRspObj;
X
Xiaoyu Wang 已提交
1596 1597
}

L
Liu Jicong 已提交
1598 1599
SMqTaosxRspObj* tmqBuildTaosxRspFromWrapper(SMqPollRspWrapper* pWrapper) {
  SMqTaosxRspObj* pRspObj = taosMemoryCalloc(1, sizeof(SMqTaosxRspObj));
1600
  pRspObj->resType = RES_TYPE__TMQ_METADATA;
L
Liu Jicong 已提交
1601 1602 1603 1604
  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;
1605
  memcpy(&pRspObj->rsp, &pWrapper->taosxRsp, sizeof(STaosxRsp));
L
Liu Jicong 已提交
1606 1607 1608

  pRspObj->resInfo.totalRows = 0;
  pRspObj->resInfo.precision = TSDB_TIME_PRECISION_MILLI;
1609
  if (!pWrapper->taosxRsp.withSchema) {
L
Liu Jicong 已提交
1610 1611 1612 1613 1614 1615
    setResSchemaInfo(&pRspObj->resInfo, pWrapper->topicHandle->schema.pSchema, pWrapper->topicHandle->schema.nCols);
  }

  return pRspObj;
}

1616 1617 1618 1619 1620 1621 1622 1623 1624 1625 1626 1627 1628 1629 1630 1631 1632 1633 1634 1635 1636 1637 1638 1639 1640 1641 1642 1643 1644 1645 1646 1647 1648
static int32_t handleErrorBeforePoll(SMqClientVg* pVg, tmq_t* pTmq) {
  atomic_store_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE);
  tsem_post(&pTmq->rspSem);
  return -1;
}

static int32_t doTmqPollImpl(tmq_t* pTmq, SMqClientTopic* pTopic, SMqClientVg* pVg, int64_t timeout) {
  SMqPollReq req = {0};
  tmqBuildConsumeReqImpl(&req, pTmq, timeout, pTopic, pVg);

  int32_t msgSize = tSerializeSMqPollReq(NULL, 0, &req);
  if (msgSize < 0) {
    return handleErrorBeforePoll(pVg, pTmq);
  }

  char* msg = taosMemoryCalloc(1, msgSize);
  if (NULL == msg) {
    return handleErrorBeforePoll(pVg, pTmq);
  }

  if (tSerializeSMqPollReq(msg, msgSize, &req) < 0) {
    taosMemoryFree(msg);
    return handleErrorBeforePoll(pVg, pTmq);
  }

  SMqPollCbParam* pParam = taosMemoryMalloc(sizeof(SMqPollCbParam));
  if (pParam == NULL) {
    taosMemoryFree(msg);
    return handleErrorBeforePoll(pVg, pTmq);
  }

  pParam->refId = pTmq->refId;
  pParam->epoch = pTmq->epoch;
X
Xiaoyu Wang 已提交
1649
  pParam->pVg = pVg;  // pVg may be released,fix it
1650 1651
  pParam->pTopic = pTopic;
  pParam->vgId = pVg->vgId;
H
Haojun Liao 已提交
1652
  pParam->requestId = req.reqId;
1653 1654 1655 1656 1657 1658 1659 1660 1661 1662 1663 1664 1665 1666 1667 1668 1669 1670 1671 1672 1673 1674 1675 1676

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

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

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

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

H
Haojun Liao 已提交
1677
  tscDebug("consumer:0x%" PRIx64 " send poll to %s vgId:%d, epoch %d, req:%s, reqId:0x%" PRIx64,
1678 1679 1680 1681 1682 1683 1684 1685 1686
           pTmq->consumerId, pTopic->topicName, pVg->vgId, pTmq->epoch, offsetFormatBuf, req.reqId);
  asyncSendMsgToServer(pTmq->pTscObj->pAppInfo->pTransporter, &pVg->epSet, &transporterId, sendInfo);

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

  return TSDB_CODE_SUCCESS;
}

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

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

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

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

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

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

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

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

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

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

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

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

1776
      int32_t consumerEpoch = atomic_load_32(&tmq->epoch);
1777 1778 1779
      SMqDataRsp* pDataRsp = &pollRspWrapper->dataRsp;

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

        // update the epset
        if (pollRspWrapper->pEpset != NULL) {
          SEp* pEp = GET_ACTIVE_EP(pollRspWrapper->pEpset);
          SEp* pOld = GET_ACTIVE_EP(&(pVg->epSet));
          tscDebug("consumer:0x%" PRIx64 " update epset vgId:%d, ep:%s:%d, old ep:%s:%d", tmq->consumerId,
                   pVg->vgId, pEp->fqdn, pEp->port, pOld->fqdn, pOld->port);
          pVg->epSet = *pollRspWrapper->pEpset;
        }

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

1795 1796 1797
        char buf[80];
        tFormatOffset(buf, 80, &pDataRsp->rspOffset);
        if (pDataRsp->blockNum == 0) {
H
Haojun Liao 已提交
1798 1799
          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);
1800
          pRspWrapper = tmqFreeRspWrapper(pRspWrapper);
L
Liu Jicong 已提交
1801
          taosFreeQitem(pollRspWrapper);
1802
        } else {  // build rsp
1803 1804
          int64_t numOfRows = 0;
          SMqRspObj* pRsp = tmqBuildRspFromWrapper(pollRspWrapper, pVg, &numOfRows);
H
Haojun Liao 已提交
1805 1806
          tmq->totalRows += numOfRows;

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

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

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

L
Liu Jicong 已提交
1844 1845 1846 1847
      if (pollRspWrapper->taosxRsp.head.epoch == consumerEpoch) {
        SMqClientVg* pVg = pollRspWrapper->vgHandle;
        pVg->currentOffset = pollRspWrapper->taosxRsp.rspOffset;
        atomic_store_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE);
H
Haojun Liao 已提交
1848

L
Liu Jicong 已提交
1849
        if (pollRspWrapper->taosxRsp.blockNum == 0) {
H
Haojun Liao 已提交
1850 1851
          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 已提交
1852
          pVg->emptyBlockReceiveTs = taosGetTimestampMs();
1853
          pRspWrapper = tmqFreeRspWrapper(pRspWrapper);
H
Haojun Liao 已提交
1854
          taosFreeQitem(pollRspWrapper);
L
Liu Jicong 已提交
1855
          continue;
H
Haojun Liao 已提交
1856 1857
        } else {
          pVg->emptyBlockReceiveTs = 0; // reset the ts
L
Liu Jicong 已提交
1858
        }
wmmhello's avatar
wmmhello 已提交
1859

L
Liu Jicong 已提交
1860
        // build rsp
wmmhello's avatar
wmmhello 已提交
1861
        void* pRsp = NULL;
1862
        int64_t numOfRows = 0;
L
Liu Jicong 已提交
1863
        if (pollRspWrapper->taosxRsp.createTableNum == 0) {
1864
          pRsp = tmqBuildRspFromWrapper(pollRspWrapper, pVg, &numOfRows);
L
Liu Jicong 已提交
1865
        } else {
wmmhello's avatar
wmmhello 已提交
1866 1867
          pRsp = tmqBuildTaosxRspFromWrapper(pollRspWrapper);
        }
H
Haojun Liao 已提交
1868

1869 1870
        tmq->totalRows += numOfRows;

H
Haojun Liao 已提交
1871 1872
        char buf[80];
        tFormatOffset(buf, 80, &pVg->currentOffset);
H
Haojun Liao 已提交
1873
        tscDebug("consumer:0x%" PRIx64 " process taosx poll rsp, vgId:%d, offset:%s, blocks:%d, rows:%" PRId64
H
Haojun Liao 已提交
1874
                 ", vg total:%" PRId64 " total:%"PRId64" reqId:0x%" PRIx64,
H
Haojun Liao 已提交
1875
                 tmq->consumerId, pVg->vgId, buf, pollRspWrapper->dataRsp.blockNum, numOfRows, pVg->numOfRows,
H
Haojun Liao 已提交
1876
                 tmq->totalRows, pollRspWrapper->reqId);
H
Haojun Liao 已提交
1877 1878

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

L
Liu Jicong 已提交
1881
      } else {
H
Haojun Liao 已提交
1882
        tscDebug("consumer:0x%" PRIx64 " vgId:%d msg discard since epoch mismatch: msg epoch %d, consumer epoch %d",
1883
                 tmq->consumerId, pollRspWrapper->vgId, pollRspWrapper->taosxRsp.head.epoch, consumerEpoch);
1884
        pRspWrapper = tmqFreeRspWrapper(pRspWrapper);
L
Liu Jicong 已提交
1885 1886
        taosFreeQitem(pollRspWrapper);
      }
X
Xiaoyu Wang 已提交
1887
    } else {
H
Haojun Liao 已提交
1888 1889
      tscDebug("consumer:0x%" PRIx64 " not data msg received", tmq->consumerId);

X
Xiaoyu Wang 已提交
1890
      bool reset = false;
1891 1892
      tmqHandleNoPollRsp(tmq, pRspWrapper, &reset);
      taosFreeQitem(pRspWrapper);
X
Xiaoyu Wang 已提交
1893
      if (pollIfReset && reset) {
1894
        tscDebug("consumer:0x%" PRIx64 ", reset and repoll", tmq->consumerId);
1895
        tmqPollImpl(tmq, timeout);
X
Xiaoyu Wang 已提交
1896 1897 1898 1899 1900
      }
    }
  }
}

1901
TAOS_RES* tmq_consumer_poll(tmq_t* tmq, int64_t timeout) {
L
Liu Jicong 已提交
1902 1903
  void*   rspObj;
  int64_t startTime = taosGetTimestampMs();
L
Liu Jicong 已提交
1904

1905
  tscDebug("consumer:0x%" PRIx64 " start to poll at %"PRId64", timeout:%" PRId64, tmq->consumerId, startTime, timeout);
L
Liu Jicong 已提交
1906

1907 1908 1909
#if 0
  tmqHandleAllDelayedTask(tmq);
  tmqPollImpl(tmq, timeout);
1910
  rspObj = tmqHandleAllRsp(tmq, timeout, false);
L
Liu Jicong 已提交
1911 1912
  if (rspObj) {
    return (TAOS_RES*)rspObj;
L
fix  
Liu Jicong 已提交
1913
  }
1914
#endif
X
Xiaoyu Wang 已提交
1915

1916
  // in no topic status, delayed task also need to be processed
L
Liu Jicong 已提交
1917
  if (atomic_load_8(&tmq->status) == TMQ_CONSUMER_STATUS__INIT) {
1918
    tscDebug("consumer:0x%" PRIx64 " poll return since consumer is init", tmq->consumerId);
1919
    taosMsleep(500);  //     sleep for a while
1920 1921 1922
    return NULL;
  }

wmmhello's avatar
wmmhello 已提交
1923
  while (atomic_load_8(&tmq->status) == TMQ_CONSUMER_STATUS__RECOVER) {
L
Liu Jicong 已提交
1924
    int32_t retryCnt = 0;
1925
    while (TSDB_CODE_MND_CONSUMER_NOT_READY == doAskEp(tmq)) {
H
Haojun Liao 已提交
1926
      if (retryCnt++ > 40) {
L
Liu Jicong 已提交
1927 1928
        return NULL;
      }
1929

H
Haojun Liao 已提交
1930
      tscDebug("consumer:0x%" PRIx64 " not ready, retry:%d/40 in 500ms", tmq->consumerId, retryCnt);
L
Liu Jicong 已提交
1931 1932 1933 1934
      taosMsleep(500);
    }
  }

X
Xiaoyu Wang 已提交
1935
  while (1) {
L
Liu Jicong 已提交
1936
    tmqHandleAllDelayedTask(tmq);
1937

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

1942
    rspObj = tmqHandleAllRsp(tmq, timeout, false);
L
Liu Jicong 已提交
1943
    if (rspObj) {
1944
      tscDebug("consumer:0x%" PRIx64 " return rsp %p", tmq->consumerId, rspObj);
L
Liu Jicong 已提交
1945
      return (TAOS_RES*)rspObj;
L
Liu Jicong 已提交
1946
    } else if (terrno == TSDB_CODE_TQ_NO_COMMITTED_OFFSET) {
1947
      tscDebug("consumer:0x%" PRIx64 " return null since no committed offset", tmq->consumerId);
L
Liu Jicong 已提交
1948
      return NULL;
X
Xiaoyu Wang 已提交
1949
    }
1950

1951
    if (timeout >= 0) {
L
Liu Jicong 已提交
1952
      int64_t currentTime = taosGetTimestampMs();
1953 1954 1955
      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 已提交
1956
                 tmq->consumerId, tmq->epoch, startTime, currentTime);
X
Xiaoyu Wang 已提交
1957 1958
        return NULL;
      }
1959
      tsem_timewait(&tmq->rspSem, (timeout - elapsedTime));
L
Liu Jicong 已提交
1960 1961
    } else {
      // use tsem_timewait instead of tsem_wait to avoid unexpected stuck
L
Liu Jicong 已提交
1962
      tsem_timewait(&tmq->rspSem, 1000);
X
Xiaoyu Wang 已提交
1963 1964 1965 1966
    }
  }
}

1967 1968 1969 1970 1971 1972 1973 1974 1975 1976 1977 1978 1979 1980
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);
1981
    }
1982
  }
1983

1984 1985
  tscDebug("consumer:0x%" PRIx64 " rows dist end", pTmq->consumerId);
}
1986

1987 1988 1989
int32_t tmq_consumer_close(tmq_t* tmq) {
  tscDebug("consumer:0x%" PRIx64" start to close consumer, status:%d", tmq->consumerId, tmq->status);
  displayConsumeStatistics(tmq);
1990

1991 1992 1993 1994 1995 1996
  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;
1997 1998 1999
      }
    }

L
Liu Jicong 已提交
2000
    int32_t     retryCnt = 0;
2001
    tmq_list_t* lst = tmq_list_new();
L
Liu Jicong 已提交
2002
    while (1) {
2003
      int32_t rsp = tmq_subscribe(tmq, lst);
L
Liu Jicong 已提交
2004 2005 2006 2007 2008 2009 2010 2011
      if (rsp != TSDB_CODE_MND_CONSUMER_NOT_READY || retryCnt > 5) {
        break;
      } else {
        retryCnt++;
        taosMsleep(500);
      }
    }

2012
    tmq_list_destroy(lst);
2013 2014
  } else {
    tscWarn("consumer:0x%" PRIx64" not in ready state, close it directly", tmq->consumerId);
L
Liu Jicong 已提交
2015
  }
H
Haojun Liao 已提交
2016

2017
  taosRemoveRef(tmqMgmt.rsetId, tmq->refId);
L
Liu Jicong 已提交
2018
  return 0;
2019
}
L
Liu Jicong 已提交
2020

L
Liu Jicong 已提交
2021 2022
const char* tmq_err2str(int32_t err) {
  if (err == 0) {
L
Liu Jicong 已提交
2023
    return "success";
L
Liu Jicong 已提交
2024
  } else if (err == -1) {
L
Liu Jicong 已提交
2025 2026 2027
    return "fail";
  } else {
    return tstrerror(err);
L
Liu Jicong 已提交
2028 2029
  }
}
L
Liu Jicong 已提交
2030

L
Liu Jicong 已提交
2031 2032 2033 2034 2035
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;
2036 2037
  } else if (TD_RES_TMQ_METADATA(res)) {
    return TMQ_RES_METADATA;
L
Liu Jicong 已提交
2038 2039 2040 2041 2042
  } else {
    return TMQ_RES_INVALID;
  }
}

L
Liu Jicong 已提交
2043
const char* tmq_get_topic_name(TAOS_RES* res) {
L
Liu Jicong 已提交
2044 2045
  if (TD_RES_TMQ(res)) {
    SMqRspObj* pRspObj = (SMqRspObj*)res;
L
Liu Jicong 已提交
2046
    return strchr(pRspObj->topic, '.') + 1;
L
Liu Jicong 已提交
2047 2048 2049
  } else if (TD_RES_TMQ_META(res)) {
    SMqMetaRspObj* pMetaRspObj = (SMqMetaRspObj*)res;
    return strchr(pMetaRspObj->topic, '.') + 1;
2050 2051 2052
  } else if (TD_RES_TMQ_METADATA(res)) {
    SMqTaosxRspObj* pRspObj = (SMqTaosxRspObj*)res;
    return strchr(pRspObj->topic, '.') + 1;
L
Liu Jicong 已提交
2053 2054 2055 2056 2057
  } else {
    return NULL;
  }
}

L
Liu Jicong 已提交
2058 2059 2060 2061
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 已提交
2062 2063 2064
  } else if (TD_RES_TMQ_META(res)) {
    SMqMetaRspObj* pMetaRspObj = (SMqMetaRspObj*)res;
    return strchr(pMetaRspObj->db, '.') + 1;
2065 2066 2067
  } else if (TD_RES_TMQ_METADATA(res)) {
    SMqTaosxRspObj* pRspObj = (SMqTaosxRspObj*)res;
    return strchr(pRspObj->db, '.') + 1;
L
Liu Jicong 已提交
2068 2069 2070 2071 2072
  } else {
    return NULL;
  }
}

L
Liu Jicong 已提交
2073 2074 2075 2076
int32_t tmq_get_vgroup_id(TAOS_RES* res) {
  if (TD_RES_TMQ(res)) {
    SMqRspObj* pRspObj = (SMqRspObj*)res;
    return pRspObj->vgId;
L
Liu Jicong 已提交
2077 2078 2079
  } else if (TD_RES_TMQ_META(res)) {
    SMqMetaRspObj* pMetaRspObj = (SMqMetaRspObj*)res;
    return pMetaRspObj->vgId;
2080
  } else if (TD_RES_TMQ_METADATA(res)) {
2081 2082
    SMqTaosxRspObj* pRspObj = (SMqTaosxRspObj*)res;
    return pRspObj->vgId;
L
Liu Jicong 已提交
2083 2084 2085 2086
  } else {
    return -1;
  }
}
L
Liu Jicong 已提交
2087 2088 2089 2090 2091 2092 2093 2094

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;
    }
2095
    return (const char*)taosArrayGetP(pRspObj->rsp.blockTbName, pRspObj->resIter);
2096 2097
  } else if (TD_RES_TMQ_METADATA(res)) {
    SMqTaosxRspObj* pRspObj = (SMqTaosxRspObj*)res;
L
Liu Jicong 已提交
2098 2099 2100
    if (!pRspObj->rsp.withTbName || pRspObj->rsp.blockTbName == NULL || pRspObj->resIter < 0 ||
        pRspObj->resIter >= pRspObj->rsp.blockNum) {
      return NULL;
2101
    }
L
Liu Jicong 已提交
2102 2103
    return (const char*)taosArrayGetP(pRspObj->rsp.blockTbName, pRspObj->resIter);
  }
L
Liu Jicong 已提交
2104 2105
  return NULL;
}
2106

2107 2108 2109 2110 2111 2112
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 已提交
2113 2114
}

2115 2116 2117 2118 2119 2120 2121 2122 2123
typedef struct SSyncCommitInfo {
  tsem_t  sem;
  int32_t code;
} SSyncCommitInfo;

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

2126 2127 2128 2129 2130 2131 2132 2133 2134 2135 2136
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) {
    asyncCommitOffset(tmq, pRes, commitCallBackFn, pInfo);
  } else {
    asyncCommitAllOffsets(tmq, commitCallBackFn, pInfo);
2137 2138
  }

2139 2140 2141 2142 2143 2144 2145 2146 2147 2148 2149 2150 2151 2152 2153 2154 2155 2156 2157 2158 2159 2160 2161 2162 2163 2164 2165 2166 2167 2168 2169 2170 2171 2172 2173 2174 2175 2176 2177 2178 2179 2180 2181 2182 2183 2184 2185 2186 2187 2188 2189 2190 2191 2192 2193 2194 2195 2196 2197
  tsem_wait(&pInfo->sem);

  code = pInfo->code;
  taosMemoryFree(pInfo);

  tscDebug("consumer:0x%"PRIx64" sync commit done, code:%s", tmq->consumerId, tstrerror(code));
  return code;
}

void updateEpCallbackFn(tmq_t* pTmq, int32_t code, SDataBuf* pDataBuf, void* param) {
  SAskEpInfo* pInfo = param;
  pInfo->code = code;

  if (code == TSDB_CODE_SUCCESS) {
    SMqRspHead* head = pDataBuf->pData;

    SMqAskEpRsp rsp;
    tDecodeSMqAskEpRsp(POINTER_SHIFT(pDataBuf->pData, sizeof(SMqRspHead)), &rsp);
    tmqUpdateEp(pTmq, head->epoch, &rsp);
    tDeleteSMqAskEpRsp(&rsp);
  }

  tsem_post(&pTmq->rspSem);
}

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

  asyncAskEp(pTmq, updateEpCallbackFn, pInfo);
  tsem_wait(&pTmq->rspSem);

  int32_t code = pInfo->code;
  taosMemoryFree(pInfo);
  return code;
}

void asyncAskEp(tmq_t* pTmq, __tmq_askep_fn_t askEpFn, void* param) {
2198
  SMqAskEpReq req = {0};
2199 2200 2201
  req.consumerId = pTmq->consumerId;
  req.epoch = pTmq->epoch;
  strcpy(req.cgroup, pTmq->groupId);
2202 2203 2204

  int32_t tlen = tSerializeSMqAskEpReq(NULL, 0, &req);
  if (tlen < 0) {
2205 2206 2207
    tscError("consumer:0x%" PRIx64 ", tSerializeSMqAskEpReq failed", pTmq->consumerId);
    askEpFn(pTmq, TSDB_CODE_INVALID_PARA, NULL, param);
    return;
2208 2209 2210 2211
  }

  void* pReq = taosMemoryCalloc(1, tlen);
  if (pReq == NULL) {
2212 2213 2214
    tscError("consumer:0x%" PRIx64 ", failed to malloc askEpReq msg, size:%d", pTmq->consumerId, tlen);
    askEpFn(pTmq, TSDB_CODE_OUT_OF_MEMORY, NULL, param);
    return;
2215 2216 2217
  }

  if (tSerializeSMqAskEpReq(pReq, tlen, &req) < 0) {
2218
    tscError("consumer:0x%" PRIx64 ", tSerializeSMqAskEpReq %d failed", pTmq->consumerId, tlen);
2219
    taosMemoryFree(pReq);
2220 2221 2222

    askEpFn(pTmq, TSDB_CODE_INVALID_PARA, NULL, param);
    return;
2223 2224 2225 2226
  }

  SMqAskEpCbParam* pParam = taosMemoryCalloc(1, sizeof(SMqAskEpCbParam));
  if (pParam == NULL) {
2227
    tscError("consumer:0x%" PRIx64 ", failed to malloc subscribe param", pTmq->consumerId);
2228
    taosMemoryFree(pReq);
2229 2230 2231

    askEpFn(pTmq, TSDB_CODE_OUT_OF_MEMORY, NULL, param);
    return;
2232 2233
  }

2234 2235 2236 2237
  pParam->refId = pTmq->refId;
  pParam->epoch = pTmq->epoch;
  pParam->pUserFn = askEpFn;
  pParam->pParam = param;
2238 2239 2240 2241 2242

  SMsgSendInfo* sendInfo = taosMemoryCalloc(1, sizeof(SMsgSendInfo));
  if (sendInfo == NULL) {
    taosMemoryFree(pParam);
    taosMemoryFree(pReq);
2243 2244
    askEpFn(pTmq, TSDB_CODE_OUT_OF_MEMORY, NULL, param);
    return;
2245 2246
  }

2247
  sendInfo->msgInfo = (SDataBuf){ .pData = pReq, .len = tlen, .handle = NULL };
2248 2249 2250 2251

  sendInfo->requestId = generateRequestId();
  sendInfo->requestObjRefId = 0;
  sendInfo->param = pParam;
2252
  sendInfo->fp = askEpCallbackFn;
2253 2254
  sendInfo->msgType = TDMT_MND_TMQ_ASK_EP;

2255 2256
  SEpSet epSet = getEpSet_s(&pTmq->pTscObj->pAppInfo->mgmtEp);
  tscDebug("consumer:0x%" PRIx64 " ask ep from mnode, reqId:0x%" PRIx64, pTmq->consumerId, sendInfo->requestId);
2257 2258

  int64_t transporterId = 0;
2259
  asyncSendMsgToServer(pTmq->pTscObj->pAppInfo->pTransporter, &epSet, &transporterId, sendInfo);
2260 2261 2262 2263 2264 2265 2266
}

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

int32_t tmqCommitDone(SMqCommitCbParamSet* pParamSet) {
2267 2268 2269
  int64_t refId = pParamSet->refId;

  tmq_t* tmq = taosAcquireRef(tmqMgmt.rsetId, refId);
2270 2271 2272 2273 2274 2275 2276
  if (tmq == NULL) {
    taosMemoryFree(pParamSet);
    terrno = TSDB_CODE_TMQ_CONSUMER_CLOSED;
    return -1;
  }

  // if no more waiting rsp
2277 2278
  pParamSet->userCb(tmq, pParamSet->code, pParamSet->userParam);
  taosMemoryFree(pParamSet);
2279 2280

  taosReleaseRef(tmqMgmt.rsetId, refId);
2281
  return 0;
2282 2283
}

2284
void commitRspCountDown(SMqCommitCbParamSet* pParamSet, int64_t consumerId, const char* pTopic, int32_t vgId) {
2285 2286
  int32_t waitingRspNum = atomic_sub_fetch_32(&pParamSet->waitingRspNum, 1);
  if (waitingRspNum == 0) {
H
Haojun Liao 已提交
2287 2288
    tscDebug("consumer:0x%" PRIx64 " topic:%s vgId:%d all commit-rsp received, commit completed", consumerId, pTopic,
             vgId);
2289
    tmqCommitDone(pParamSet);
H
Haojun Liao 已提交
2290 2291 2292
  } else {
    tscDebug("consumer:0x%" PRIx64 " topic:%s vgId:%d commit-rsp received, remain:%d", consumerId, pTopic, vgId,
             waitingRspNum);
2293 2294
  }
}