clientTmq.c 72.4 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
  if (strcasecmp(key, "enable.heartbeat.background") == 0) {
327 328 329 330 331 332 333 334
//    if (strcasecmp(value, "true") == 0) {
//      conf->hbBgEnable = true;
//      return TMQ_CONF_OK;
//    } else if (strcasecmp(value, "false") == 0) {
//      conf->hbBgEnable = false;
//      return TMQ_CONF_OK;
//    } else {
      tscError("the default value of enable.heartbeat.background is true, can not be seted");
L
Liu Jicong 已提交
335
      return TMQ_CONF_INVALID;
336
//    }
L
Liu Jicong 已提交
337 338
  }

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

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

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

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

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

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

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

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

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

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

396 397 398 399 400
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;

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

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

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

420 421 422
// Two problems do not need to be addressed here
// 1. update to of epset. the response of poll request will automatically handle this problem
// 2. commit failure. This one needs to be resolved.
H
Haojun Liao 已提交
423
static int32_t tmqCommitCb(void* param, SDataBuf* pBuf, int32_t code) {
424
  SMqCommitCbParam*    pParam = (SMqCommitCbParam*)param;
425
  SMqCommitCbParamSet* pParamSet = (SMqCommitCbParamSet*)pParam->params;
H
Haojun Liao 已提交
426

427 428 429 430 431 432 433 434 435 436 437 438 439 440 441 442 443 444 445 446 447 448 449 450
//  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);
//
451
//    commitRspCountDown(pParamSet, pParam->pTmq->consumerId, pParam->topicName, pParam->vgId);
452 453 454 455
//    return 0;
//  }
//
//  // todo replace the pTmq with refId
456

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

557 558 559 560 561 562 563 564 565 566 567 568 569
static void asyncCommitOffset(tmq_t* tmq, const TAOS_RES* pRes, tmq_commit_cb* pCommitFp, void* userParam) {
  char*   pTopicName = NULL;
  int32_t vgId = 0;
  int32_t code = 0;

  if (pRes == NULL || tmq == NULL) {
    pCommitFp(tmq, TSDB_CODE_INVALID_PARA, userParam);
    return;
  }

  if (TD_RES_TMQ(pRes)) {
    SMqRspObj* pRspObj = (SMqRspObj*)pRes;
    pTopicName = pRspObj->topic;
L
Liu Jicong 已提交
570
    vgId = pRspObj->vgId;
571 572 573
  } else if (TD_RES_TMQ_META(pRes)) {
    SMqMetaRspObj* pMetaRspObj = (SMqMetaRspObj*)pRes;
    pTopicName = pMetaRspObj->topic;
L
Liu Jicong 已提交
574
    vgId = pMetaRspObj->vgId;
575 576 577
  } else if (TD_RES_TMQ_METADATA(pRes)) {
    SMqTaosxRspObj* pRspObj = (SMqTaosxRspObj*)pRes;
    pTopicName = pRspObj->topic;
578
    vgId = pRspObj->vgId;
L
Liu Jicong 已提交
579
  } else {
580 581
    pCommitFp(tmq, TSDB_CODE_TMQ_INVALID_MSG, userParam);
    return;
L
Liu Jicong 已提交
582 583 584 585
  }

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

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

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

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

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

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

615 616 617 618 619 620 621 622
  SMqClientTopic* pTopic = taosArrayGet(tmq->clientTopics, i);

  int32_t j = 0;
  int32_t numOfVgroups = taosArrayGetSize(pTopic->vgs);
  for (j = 0; j < numOfVgroups; j++) {
    SMqClientVg* pVg = taosArrayGet(pTopic->vgs, j);
    if (pVg->vgId == vgId) {
      break;
L
Liu Jicong 已提交
623
    }
L
Liu Jicong 已提交
624
  }
L
Liu Jicong 已提交
625

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

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

638 639 640 641 642 643 644 645
    // 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 已提交
646 647 648
  }
}

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

L
Liu Jicong 已提交
863 864 865 866 867 868
    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;
869 870
    taosMemoryFreeClear(pRsp->pEpset);

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

L
Liu Jicong 已提交
876 877 878 879 880 881 882 883
    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);
  }
884 885

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

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

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

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

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

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

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

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

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

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

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

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

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

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

982 983 984 985 986 987 988 989 990
static void tmqMgmtInit(void) {
  tmqInitRes = 0;
  tmqMgmt.timer = taosTmrInit(1000, 100, 360000, "TMQ");

  if (tmqMgmt.timer == NULL) {
    tmqInitRes = TSDB_CODE_OUT_OF_MEMORY;
  }

  tmqMgmt.rsetId = taosOpenRef(10000, tmqFreeImpl);
991
  if (tmqMgmt.rsetId < 0) {
992 993 994 995
    tmqInitRes = terrno;
  }
}

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

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

L
Liu Jicong 已提交
1010 1011 1012
  const char* user = conf->user == NULL ? TSDB_DEFAULT_USER : conf->user;
  const char* pass = conf->pass == NULL ? TSDB_DEFAULT_PASS : conf->pass;

L
Liu Jicong 已提交
1013 1014 1015
  pTmq->clientTopics = taosArrayInit(0, sizeof(SMqClientTopic));
  pTmq->mqueue = taosOpenQueue();
  pTmq->delayedTask = taosOpenQueue();
H
Haojun Liao 已提交
1016
  pTmq->qall = taosAllocateQall();
L
Liu Jicong 已提交
1017

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

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

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

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

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

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

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

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

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

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

1083
  return pTmq;
1084

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

L
Liu Jicong 已提交
1243
  if (code != 0) {
H
Haojun Liao 已提交
1244 1245
    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 已提交
1246

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

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

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

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

X
Xiaoyu Wang 已提交
1270 1271 1272
  int32_t msgEpoch = ((SMqRspHead*)pMsg->pData)->epoch;
  int32_t tmqEpoch = atomic_load_32(&tmq->epoch);
  if (msgEpoch < tmqEpoch) {
L
Liu Jicong 已提交
1273
    // do not write into queue since updating epoch reset
H
Haojun Liao 已提交
1274 1275 1276
    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);

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

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

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

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

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

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

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

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

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

H
Haojun Liao 已提交
1339 1340 1341
  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);

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

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

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

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

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

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

H
Haojun Liao 已提交
1363 1364 1365 1366 1367 1368 1369 1370 1371 1372 1373 1374 1375 1376 1377 1378 1379
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 已提交
1380 1381

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

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

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

    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) {
1416 1417
  bool set = false;

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

1519 1520 1521 1522 1523 1524 1525 1526
    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 已提交
1527
  } else {
1528 1529 1530
    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);
1531
  }
L
Liu Jicong 已提交
1532

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

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

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

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

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

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

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

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

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

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

1588 1589 1590 1591 1592
  // 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;
1593
    (*numOfRows) += rows;
1594 1595
  }

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

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

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

  return pRspObj;
}

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

  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 已提交
1678
  tscDebug("consumer:0x%" PRIx64 " send poll to %s vgId:%d, epoch %d, req:%s, reqId:0x%" PRIx64,
1679 1680 1681 1682 1683 1684 1685 1686 1687
           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;
}

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

L
Liu Jicong 已提交
1845 1846 1847 1848
      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 已提交
1849

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

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

1870 1871
        tmq->totalRows += numOfRows;

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

2116 2117 2118 2119 2120 2121 2122 2123 2124
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);
2125
}
2126

2127 2128 2129 2130 2131 2132 2133 2134 2135
int32_t tmq_commit_sync(tmq_t* tmq, const TAOS_RES* pRes) {
  int32_t code = 0;

  SSyncCommitInfo* pInfo = taosMemoryMalloc(sizeof(SSyncCommitInfo));
  tsem_init(&pInfo->sem, 0, 0);
  pInfo->code = 0;

  if (pRes == NULL) {
    asyncCommitAllOffsets(tmq, commitCallBackFn, pInfo);
H
Haojun Liao 已提交
2136 2137
  } else {
    asyncCommitOffset(tmq, pRes, commitCallBackFn, pInfo);
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 2198
  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) {
2199
  SMqAskEpReq req = {0};
2200 2201 2202
  req.consumerId = pTmq->consumerId;
  req.epoch = pTmq->epoch;
  strcpy(req.cgroup, pTmq->groupId);
2203 2204 2205

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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