clientTmq.c 70.7 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  100
28
#define DEFAULT_AUTO_COMMIT_INTERVAL    5000
29

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

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

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

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

L
Liu Jicong 已提交
51
struct tmq_list_t {
L
Liu Jicong 已提交
52
  SArray container;
L
Liu Jicong 已提交
53
};
L
Liu Jicong 已提交
54

L
Liu Jicong 已提交
55
struct tmq_conf_t {
56 57 58 59 60 61 62 63
  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;
64 65 66 67 68
  uint16_t       port;
  int32_t        autoCommitInterval;
  char*          ip;
  char*          user;
  char*          pass;
69
  tmq_commit_cb* commitCb;
L
Liu Jicong 已提交
70
  void*          commitCbUserParam;
L
Liu Jicong 已提交
71 72 73
};

struct tmq_t {
74 75 76 77 78 79 80 81 82 83
  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 已提交
84 85
  tmq_commit_cb* commitCb;
  void*          commitCbUserParam;
L
Liu Jicong 已提交
86 87 88 89

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

L
Liu Jicong 已提交
98
  // timer
99 100 101 102
  tmr_h         hbLiveTimer;
  tmr_h         epTimer;
  tmr_h         reportTimer;
  tmr_h         commitTimer;
H
Haojun Liao 已提交
103 104 105 106 107 108 109
  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 已提交
110 111
};

X
Xiaoyu Wang 已提交
112 113 114 115 116 117 118 119
enum {
  TMQ_VG_STATUS__IDLE = 0,
  TMQ_VG_STATUS__WAIT,
};

enum {
  TMQ_CONSUMER_STATUS__INIT = 0,
  TMQ_CONSUMER_STATUS__READY,
120
  TMQ_CONSUMER_STATUS__NO_TOPIC,
L
Liu Jicong 已提交
121
  TMQ_CONSUMER_STATUS__RECOVER,
L
Liu Jicong 已提交
122 123
};

L
Liu Jicong 已提交
124
enum {
125
  TMQ_DELAYED_TASK__ASK_EP = 1,
L
Liu Jicong 已提交
126 127 128 129
  TMQ_DELAYED_TASK__REPORT,
  TMQ_DELAYED_TASK__COMMIT,
};

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

L
Liu Jicong 已提交
142
typedef struct {
143 144 145
  char           topicName[TSDB_TOPIC_FNAME_LEN];
  char           db[TSDB_DB_FNAME_LEN];
  SArray*        vgs;  // SArray<SMqClientVg>
L
Liu Jicong 已提交
146
  SSchemaWrapper schema;
147 148
} SMqClientTopic;

L
Liu Jicong 已提交
149 150 151 152 153
typedef struct {
  int8_t          tmqRspType;
  int32_t         epoch;
  SMqClientVg*    vgHandle;
  SMqClientTopic* topicHandle;
H
Haojun Liao 已提交
154
  uint64_t        reqId;
155
  SEpSet*         pEpset;
L
Liu Jicong 已提交
156
  union {
L
Liu Jicong 已提交
157 158
    SMqDataRsp dataRsp;
    SMqMetaRsp metaRsp;
L
Liu Jicong 已提交
159
    STaosxRsp  taosxRsp;
L
Liu Jicong 已提交
160
  };
L
Liu Jicong 已提交
161 162
} SMqPollRspWrapper;

L
Liu Jicong 已提交
163
typedef struct {
164 165
  int64_t refId;
  int32_t epoch;
L
Liu Jicong 已提交
166 167
  tsem_t  rspSem;
  int32_t rspErr;
L
Liu Jicong 已提交
168
} SMqSubscribeCbParam;
L
Liu Jicong 已提交
169

L
Liu Jicong 已提交
170
typedef struct {
171 172
  int64_t refId;
  int32_t epoch;
L
Liu Jicong 已提交
173
  int32_t code;
L
Liu Jicong 已提交
174
  int32_t async;
X
Xiaoyu Wang 已提交
175
  tsem_t  rspSem;
176 177
} SMqAskEpCbParam;

L
Liu Jicong 已提交
178
typedef struct {
179 180
  int64_t         refId;
  int32_t         epoch;
L
Liu Jicong 已提交
181
  SMqClientVg*    pVg;
L
Liu Jicong 已提交
182
  SMqClientTopic* pTopic;
L
Liu Jicong 已提交
183
  int32_t         vgId;
L
Liu Jicong 已提交
184
  tsem_t          rspSem;
H
Haojun Liao 已提交
185
  uint64_t        requestId; // request id for debug purpose
X
Xiaoyu Wang 已提交
186
} SMqPollCbParam;
187

188
typedef struct {
189 190
  int64_t        refId;
  int32_t        epoch;
L
Liu Jicong 已提交
191 192
  int8_t         automatic;
  int8_t         async;
L
Liu Jicong 已提交
193 194
  int32_t        waitingRspNum;
  int32_t        totalRspNum;
L
Liu Jicong 已提交
195
  int32_t        rspErr;
196
  tmq_commit_cb* userCb;
L
Liu Jicong 已提交
197 198 199 200
  /*SArray*        successfulOffsets;*/
  /*SArray*        failedOffsets;*/
  void*  userParam;
  tsem_t rspSem;
201 202 203 204 205
} SMqCommitCbParamSet;

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

211
static int32_t tmqAskEp(tmq_t* tmq, bool async);
212 213
static int32_t makeTopicVgroupKey(char* dst, const char* topicName, int32_t vg);
static int32_t tmqCommitDone(SMqCommitCbParamSet* pParamSet);
214 215
static int32_t doSendCommitMsg(tmq_t* tmq, SMqClientVg* pVg, const char* pTopicName, SMqCommitCbParamSet* pParamSet,
                               int32_t index, int32_t totalVgroups);
216
static void tmqCommitRspCountDown(SMqCommitCbParamSet* pParamSet, int64_t consumerId, const char* pTopic, int32_t vgId);
217

218
tmq_conf_t* tmq_conf_new() {
wafwerar's avatar
wafwerar 已提交
219
  tmq_conf_t* conf = taosMemoryCalloc(1, sizeof(tmq_conf_t));
220 221 222 223 224
  if (conf == NULL) {
    terrno = TSDB_CODE_OUT_OF_MEMORY;
    return conf;
  }

225
  conf->withTbName = false;
L
Liu Jicong 已提交
226
  conf->autoCommit = true;
227
  conf->autoCommitInterval = DEFAULT_AUTO_COMMIT_INTERVAL;
228
  conf->resetOffset = TMQ_OFFSET__RESET_EARLIEAST;
229
  conf->hbBgEnable = true;
230

231 232 233
  return conf;
}

L
Liu Jicong 已提交
234
void tmq_conf_destroy(tmq_conf_t* conf) {
L
Liu Jicong 已提交
235
  if (conf) {
236 237 238 239 240 241 242 243 244
    if (conf->ip) {
      taosMemoryFree(conf->ip);
    }
    if (conf->user) {
      taosMemoryFree(conf->user);
    }
    if (conf->pass) {
      taosMemoryFree(conf->pass);
    }
L
Liu Jicong 已提交
245 246
    taosMemoryFree(conf);
  }
L
Liu Jicong 已提交
247 248 249
}

tmq_conf_res_t tmq_conf_set(tmq_conf_t* conf, const char* key, const char* value) {
250
  if (strcasecmp(key, "group.id") == 0) {
L
Liu Jicong 已提交
251
    tstrncpy(conf->groupId, value, TSDB_CGROUP_LEN);
L
Liu Jicong 已提交
252
    return TMQ_CONF_OK;
253
  }
L
Liu Jicong 已提交
254

255
  if (strcasecmp(key, "client.id") == 0) {
L
Liu Jicong 已提交
256
    tstrncpy(conf->clientId, value, 256);
L
Liu Jicong 已提交
257 258
    return TMQ_CONF_OK;
  }
L
Liu Jicong 已提交
259

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

272
  if (strcasecmp(key, "auto.commit.interval.ms") == 0) {
273
    conf->autoCommitInterval = taosStr2int64(value);
L
Liu Jicong 已提交
274 275 276
    return TMQ_CONF_OK;
  }

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

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

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

316
  if (strcasecmp(key, "experimental.snapshot.batch.size") == 0) {
317
    conf->snapBatchSize = taosStr2int64(value);
L
Liu Jicong 已提交
318 319 320
    return TMQ_CONF_OK;
  }

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

333
  if (strcasecmp(key, "td.connect.ip") == 0) {
334
    conf->ip = taosStrdup(value);
L
Liu Jicong 已提交
335 336
    return TMQ_CONF_OK;
  }
337

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

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

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

353
  if (strcasecmp(key, "td.connect.db") == 0) {
L
Liu Jicong 已提交
354 355 356
    return TMQ_CONF_OK;
  }

L
Liu Jicong 已提交
357
  return TMQ_CONF_UNKNOWN;
358 359 360
}

tmq_list_t* tmq_list_new() {
L
Liu Jicong 已提交
361
  return (tmq_list_t*)taosArrayInit(0, sizeof(void*));
362 363
}

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

L
Liu Jicong 已提交
375
void tmq_list_destroy(tmq_list_t* list) {
L
Liu Jicong 已提交
376
  SArray* container = &list->container;
L
Liu Jicong 已提交
377
  taosArrayDestroyP(container, taosMemoryFree);
L
Liu Jicong 已提交
378 379
}

L
Liu Jicong 已提交
380 381 382 383 384 385 386 387 388 389
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;
}

390 391 392 393 394
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;

395
  for(int32_t i = 0; i < numOfTopics; ++i) {
396 397
    SMqClientTopic* pTopic = taosArrayGet(pTopicList, i);
    if (strcmp(pTopic->topicName, pName) != 0) {
398 399 400
      continue;
    }

401 402
    *numOfVgroups = taosArrayGetSize(pTopic->vgs);
    for (int32_t j = 0; j < (*numOfVgroups); ++j) {
403
      SMqClientVg* pClientVg = taosArrayGet(pTopic->vgs, j);
404 405 406
      if (pClientVg->vgId == vgId) {
        *index = j;
        return pClientVg;
407 408
      }
    }
L
Liu Jicong 已提交
409
  }
410 411

  return NULL;
L
Liu Jicong 已提交
412
}
413

414 415 416
// 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 已提交
417
static int32_t tmqCommitCb(void* param, SDataBuf* pBuf, int32_t code) {
418
  SMqCommitCbParam*    pParam = (SMqCommitCbParam*)param;
419
  SMqCommitCbParamSet* pParamSet = (SMqCommitCbParamSet*)pParam->params;
H
Haojun Liao 已提交
420

421 422 423 424 425 426 427 428 429 430 431 432 433 434 435 436 437 438 439 440 441 442 443 444 445 446 447 448 449
//  if (code != TSDB_CODE_SUCCESS) { // if commit offset failed, let's try again
//    taosThreadMutexLock(&pParam->pTmq->lock);
//    int32_t numOfVgroups, index;
//    SMqClientVg* pVg = foundClientVg(pParam->pTmq->clientTopics, pParam->topicName, pParam->vgId, &index, &numOfVgroups);
//    if (pVg == NULL) {
//      tscDebug("consumer:0x%" PRIx64
//               " subKey:%s vgId:%d commit failed, code:%s has been transferred to other consumer, no need retry ordinal:%d/%d",
//               pParam->pTmq->consumerId, pParam->pOffset->subKey, pParam->vgId, tstrerror(code), index + 1, numOfVgroups);
//    } else { // let's retry the commit
//      int32_t code1 = doSendCommitMsg(pParam->pTmq, pVg, pParam->topicName, pParamSet, index, numOfVgroups);
//      if (code1 != TSDB_CODE_SUCCESS) {  // retry failed.
//        tscError("consumer:0x%" PRIx64 " topic:%s vgId:%d offset:%" PRId64
//                 " retry failed, ignore this commit. code:%s ordinal:%d/%d",
//                 pParam->pTmq->consumerId, pParam->topicName, pVg->vgId, pVg->committedOffset.version,
//                 tstrerror(terrno), index + 1, numOfVgroups);
//      }
//    }
//
//    taosThreadMutexUnlock(&pParam->pTmq->lock);
//
//    taosMemoryFree(pParam->pOffset);
//    taosMemoryFree(pBuf->pData);
//    taosMemoryFree(pBuf->pEpSet);
//
//    tmqCommitRspCountDown(pParamSet, pParam->pTmq->consumerId, pParam->topicName, pParam->vgId);
//    return 0;
//  }
//
//  // todo replace the pTmq with refId
450

L
Liu Jicong 已提交
451
  taosMemoryFree(pParam->pOffset);
L
Liu Jicong 已提交
452
  taosMemoryFree(pBuf->pData);
dengyihao's avatar
dengyihao 已提交
453
  taosMemoryFree(pBuf->pEpSet);
L
Liu Jicong 已提交
454

455
  tmqCommitRspCountDown(pParamSet, pParam->pTmq->consumerId, pParam->topicName, pParam->vgId);
456 457 458
  return 0;
}

459 460
static int32_t doSendCommitMsg(tmq_t* tmq, SMqClientVg* pVg, const char* pTopicName, SMqCommitCbParamSet* pParamSet,
                               int32_t index, int32_t totalVgroups) {
L
Liu Jicong 已提交
461 462 463 464 465
  STqOffset* pOffset = taosMemoryCalloc(1, sizeof(STqOffset));
  if (pOffset == NULL) {
    terrno = TSDB_CODE_OUT_OF_MEMORY;
    return -1;
  }
466

L
Liu Jicong 已提交
467
  pOffset->val = pVg->currentOffset;
468

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

474 475
  int32_t len = 0;
  int32_t code = 0;
L
Liu Jicong 已提交
476 477 478 479
  tEncodeSize(tEncodeSTqOffset, pOffset, len, code);
  if (code < 0) {
    return -1;
  }
480

L
Liu Jicong 已提交
481
  void* buf = taosMemoryCalloc(1, sizeof(SMsgHead) + len);
L
Liu Jicong 已提交
482 483 484 485
  if (buf == NULL) {
    taosMemoryFree(pOffset);
    return -1;
  }
486

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

L
Liu Jicong 已提交
489 490 491 492 493
  void* abuf = POINTER_SHIFT(buf, sizeof(SMsgHead));

  SEncoder encoder;
  tEncoderInit(&encoder, abuf, len);
  tEncodeSTqOffset(&encoder, pOffset);
L
Liu Jicong 已提交
494
  tEncoderClear(&encoder);
L
Liu Jicong 已提交
495 496

  // build param
497
  SMqCommitCbParam* pParam = taosMemoryCalloc(1, sizeof(SMqCommitCbParam));
L
Liu Jicong 已提交
498
  if (pParam == NULL) {
L
Liu Jicong 已提交
499
    taosMemoryFree(pOffset);
L
Liu Jicong 已提交
500 501 502
    taosMemoryFree(buf);
    return -1;
  }
503

L
Liu Jicong 已提交
504 505
  pParam->params = pParamSet;
  pParam->pOffset = pOffset;
H
Haojun Liao 已提交
506 507 508
  pParam->vgId = pVg->vgId;
  pParam->pTmq = tmq;

H
Haojun Liao 已提交
509
  tstrncpy(pParam->topicName, pTopicName, tListLen(pParam->topicName));
L
Liu Jicong 已提交
510 511 512 513

  // build send info
  SMsgSendInfo* pMsgSendInfo = taosMemoryCalloc(1, sizeof(SMsgSendInfo));
  if (pMsgSendInfo == NULL) {
L
Liu Jicong 已提交
514
    taosMemoryFree(pOffset);
L
Liu Jicong 已提交
515 516
    taosMemoryFree(buf);
    taosMemoryFree(pParam);
L
Liu Jicong 已提交
517 518
    return -1;
  }
519

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

  pMsgSendInfo->requestId = generateRequestId();
  pMsgSendInfo->requestObjRefId = 0;
  pMsgSendInfo->param = pParam;
L
Liu Jicong 已提交
529
  pMsgSendInfo->paramFreeFp = taosMemoryFree;
530
  pMsgSendInfo->fp = tmqCommitCb;
L
Liu Jicong 已提交
531
  pMsgSendInfo->msgType = TDMT_VND_TMQ_COMMIT_OFFSET;
L
Liu Jicong 已提交
532

L
Liu Jicong 已提交
533 534 535
  atomic_add_fetch_32(&pParamSet->waitingRspNum, 1);
  atomic_add_fetch_32(&pParamSet->totalRspNum, 1);

H
Haojun Liao 已提交
536 537 538 539 540 541
  SEp* pEp = GET_ACTIVE_EP(&pVg->epSet);
  tscDebug("consumer:0x%" PRIx64 " topic:%s on vgId:%d send offset:%" PRId64 " prev:%" PRId64
           ", ep:%s:%d, ordinal:%d/%d, req:0x%" PRIx64,
           tmq->consumerId, pOffset->subKey, pVg->vgId, pOffset->val.version, pVg->committedOffset.version, pEp->fqdn,
           pEp->port, index + 1, totalVgroups, pMsgSendInfo->requestId);

L
Liu Jicong 已提交
542 543 544 545 546
  int64_t transporterId = 0;
  asyncSendMsgToServer(tmq->pTscObj->pAppInfo->pTransporter, &pVg->epSet, &transporterId, pMsgSendInfo);
  return 0;
}

H
Haojun Liao 已提交
547
static int32_t tmqCommitMsgImpl(tmq_t* tmq, const TAOS_RES* msg, int8_t async, tmq_commit_cb* userCb, void* userParam) {
L
Liu Jicong 已提交
548 549 550 551 552 553 554 555 556 557
  char*   topic;
  int32_t vgId;
  if (TD_RES_TMQ(msg)) {
    SMqRspObj* pRspObj = (SMqRspObj*)msg;
    topic = pRspObj->topic;
    vgId = pRspObj->vgId;
  } else if (TD_RES_TMQ_META(msg)) {
    SMqMetaRspObj* pMetaRspObj = (SMqMetaRspObj*)msg;
    topic = pMetaRspObj->topic;
    vgId = pMetaRspObj->vgId;
L
Liu Jicong 已提交
558
  } else if (TD_RES_TMQ_METADATA(msg)) {
559 560 561
    SMqTaosxRspObj* pRspObj = (SMqTaosxRspObj*)msg;
    topic = pRspObj->topic;
    vgId = pRspObj->vgId;
L
Liu Jicong 已提交
562 563 564 565 566 567 568 569 570
  } else {
    return TSDB_CODE_TMQ_INVALID_MSG;
  }

  SMqCommitCbParamSet* pParamSet = taosMemoryCalloc(1, sizeof(SMqCommitCbParamSet));
  if (pParamSet == NULL) {
    terrno = TSDB_CODE_OUT_OF_MEMORY;
    return -1;
  }
H
Haojun Liao 已提交
571

572 573
  pParamSet->refId = tmq->refId;
  pParamSet->epoch = tmq->epoch;
L
Liu Jicong 已提交
574 575 576 577 578 579
  pParamSet->automatic = 0;
  pParamSet->async = async;
  pParamSet->userCb = userCb;
  pParamSet->userParam = userParam;
  tsem_init(&pParamSet->rspSem, 0, 0);

L
Liu Jicong 已提交
580 581
  int32_t code = -1;

H
Haojun Liao 已提交
582
  taosThreadMutexLock(&tmq->lock);
L
Liu Jicong 已提交
583 584
  for (int32_t i = 0; i < taosArrayGetSize(tmq->clientTopics); i++) {
    SMqClientTopic* pTopic = taosArrayGet(tmq->clientTopics, i);
585 586 587 588 589 590
    if (strcmp(pTopic->topicName, topic) != 0) {
      continue;
    }

    int32_t numOfVgroups = taosArrayGetSize(pTopic->vgs);
    for (int32_t j = 0; j < numOfVgroups; j++) {
L
Liu Jicong 已提交
591
      SMqClientVg* pVg = taosArrayGet(pTopic->vgs, j);
H
Haojun Liao 已提交
592 593 594
      if (pVg->vgId != vgId) {
        continue;
      }
L
Liu Jicong 已提交
595

L
Liu Jicong 已提交
596
      if (pVg->currentOffset.type > 0 && !tOffsetEqual(&pVg->currentOffset, &pVg->committedOffset)) {
597
        if (doSendCommitMsg(tmq, pVg, pTopic->topicName, pParamSet, j, numOfVgroups) < 0) {
L
Liu Jicong 已提交
598 599
          tsem_destroy(&pParamSet->rspSem);
          taosMemoryFree(pParamSet);
L
Liu Jicong 已提交
600
          goto FAIL;
L
Liu Jicong 已提交
601
        }
L
Liu Jicong 已提交
602
        goto HANDLE_RSP;
L
Liu Jicong 已提交
603 604
      }
    }
L
Liu Jicong 已提交
605
  }
L
Liu Jicong 已提交
606

L
Liu Jicong 已提交
607 608 609 610
HANDLE_RSP:
  if (pParamSet->totalRspNum == 0) {
    tsem_destroy(&pParamSet->rspSem);
    taosMemoryFree(pParamSet);
H
Haojun Liao 已提交
611
    taosThreadMutexUnlock(&tmq->lock);
L
Liu Jicong 已提交
612 613 614
    return 0;
  }

L
Liu Jicong 已提交
615
  if (!async) {
H
Haojun Liao 已提交
616
    taosThreadMutexUnlock(&tmq->lock);
L
Liu Jicong 已提交
617 618 619
    tsem_wait(&pParamSet->rspSem);
    code = pParamSet->rspErr;
    tsem_destroy(&pParamSet->rspSem);
L
Liu Jicong 已提交
620
    taosMemoryFree(pParamSet);
L
Liu Jicong 已提交
621 622 623 624 625 626
    return code;
  } else {
    code = 0;
  }

FAIL:
H
Haojun Liao 已提交
627
  taosThreadMutexUnlock(&tmq->lock);
L
Liu Jicong 已提交
628 629 630
  if (code != 0 && async) {
    userCb(tmq, code, userParam);
  }
H
Haojun Liao 已提交
631

L
Liu Jicong 已提交
632 633 634
  return 0;
}

635
static int32_t doAutoCommit(tmq_t* tmq, int8_t automatic, int8_t async, tmq_commit_cb* userCb, void* userParam) {
L
Liu Jicong 已提交
636 637
  int32_t code = -1;

638 639
  SMqCommitCbParamSet* pParamSet = taosMemoryCalloc(1, sizeof(SMqCommitCbParamSet));
  if (pParamSet == NULL) {
L
Liu Jicong 已提交
640 641 642 643 644 645 646 647
    code = TSDB_CODE_OUT_OF_MEMORY;
    if (async) {
      if (automatic) {
        tmq->commitCb(tmq, code, tmq->commitCbUserParam);
      } else {
        userCb(tmq, code, userParam);
      }
    }
648 649
    return -1;
  }
650 651 652 653

  pParamSet->refId = tmq->refId;
  pParamSet->epoch = tmq->epoch;

654 655 656 657 658 659
  pParamSet->automatic = automatic;
  pParamSet->async = async;
  pParamSet->userCb = userCb;
  pParamSet->userParam = userParam;
  tsem_init(&pParamSet->rspSem, 0, 0);

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

H
Haojun Liao 已提交
663
  taosThreadMutexLock(&tmq->lock);
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 677 678 679 680 681
      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);
        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 689
        tscDebug("consumer:0x%" PRIx64 " topic:%s vgId:%d, no commit, current:%" PRId64 ", ordinal:%d/%d",
                 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 697
  taosThreadMutexUnlock(&tmq->lock);

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

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

708 709 710 711
  if (!async) {
    tsem_wait(&pParamSet->rspSem);
    code = pParamSet->rspErr;
    tsem_destroy(&pParamSet->rspSem);
712
    taosMemoryFree(pParamSet);
L
Liu Jicong 已提交
713
#if 0
714 715
    taosArrayDestroyP(pParamSet->successfulOffsets, taosMemoryFree);
    taosArrayDestroyP(pParamSet->failedOffsets, taosMemoryFree);
L
Liu Jicong 已提交
716
#endif
L
Liu Jicong 已提交
717
  }
718

L
Liu Jicong 已提交
719 720 721
  return code;
}

722 723
static int32_t tmqCommitInner(tmq_t* tmq, const TAOS_RES* msg, int8_t automatic, int8_t async, tmq_commit_cb* userCb,
                              void* userParam) {
H
Haojun Liao 已提交
724
  if (msg) { // user invoked commit
L
Liu Jicong 已提交
725
    return tmqCommitMsgImpl(tmq, msg, async, userCb, userParam);
726
  } else {  // this for auto commit
727
    return doAutoCommit(tmq, automatic, async, userCb, userParam);
L
Liu Jicong 已提交
728
  }
729 730
}

731 732
static void generateTimedTask(int64_t refId, int32_t type) {
  tmq_t* tmq = taosAcquireRef(tmqMgmt.rsetId, refId);
733
  if (tmq != NULL) {
S
Shengliang Guan 已提交
734
    int8_t* pTaskType = taosAllocateQitem(sizeof(int8_t), DEF_QITEM, 0);
735
    *pTaskType = type;
736 737 738
    taosWriteQitem(tmq->delayedTask, pTaskType);
    tsem_post(&tmq->rspSem);
  }
739
  taosReleaseRef(tmqMgmt.rsetId, refId);
740 741 742 743 744
}

void tmqAssignAskEpTask(void* param, void* tmrId) {
  int64_t refId = *(int64_t*)param;
  generateTimedTask(refId, TMQ_DELAYED_TASK__ASK_EP);
745
  taosMemoryFree(param);
L
Liu Jicong 已提交
746 747 748
}

void tmqAssignDelayedCommitTask(void* param, void* tmrId) {
749
  int64_t refId = *(int64_t*)param;
750
  generateTimedTask(refId, TMQ_DELAYED_TASK__COMMIT);
751
  taosMemoryFree(param);
L
Liu Jicong 已提交
752 753 754
}

void tmqAssignDelayedReportTask(void* param, void* tmrId) {
755 756 757
  int64_t refId = *(int64_t*)param;
  tmq_t*  tmq = taosAcquireRef(tmqMgmt.rsetId, refId);
  if (tmq != NULL) {
S
Shengliang Guan 已提交
758
    int8_t* pTaskType = taosAllocateQitem(sizeof(int8_t), DEF_QITEM, 0);
759 760 761 762
    *pTaskType = TMQ_DELAYED_TASK__REPORT;
    taosWriteQitem(tmq->delayedTask, pTaskType);
    tsem_post(&tmq->rspSem);
  }
763 764

  taosReleaseRef(tmqMgmt.rsetId, refId);
765
  taosMemoryFree(param);
L
Liu Jicong 已提交
766 767
}

768
int32_t tmqHbCb(void* param, SDataBuf* pMsg, int32_t code) {
dengyihao's avatar
dengyihao 已提交
769 770 771 772
  if (pMsg) {
    taosMemoryFree(pMsg->pData);
    taosMemoryFree(pMsg->pEpSet);
  }
773 774 775 776
  return 0;
}

void tmqSendHbReq(void* param, void* tmrId) {
777
  int64_t refId = *(int64_t*)param;
778

779 780
  tmq_t*  tmq = taosAcquireRef(tmqMgmt.rsetId, refId);
  if (tmq == NULL) {
L
Liu Jicong 已提交
781
    taosMemoryFree(param);
782 783
    return;
  }
D
dapan1121 已提交
784 785 786 787 788

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

L
Liu Jicong 已提交
789
  int32_t tlen = tSerializeSMqHbReq(NULL, 0, &req);
D
dapan1121 已提交
790 791
  if (tlen < 0) {
    tscError("tSerializeSMqHbReq failed");
792
    goto OVER;
D
dapan1121 已提交
793
  }
794

L
Liu Jicong 已提交
795
  void* pReq = taosMemoryCalloc(1, tlen);
D
dapan1121 已提交
796 797
  if (tlen < 0) {
    tscError("failed to malloc MqHbReq msg, size:%d", tlen);
798
    goto OVER;
D
dapan1121 已提交
799
  }
800

D
dapan1121 已提交
801 802 803
  if (tSerializeSMqHbReq(pReq, tlen, &req) < 0) {
    tscError("tSerializeSMqHbReq %d failed", tlen);
    taosMemoryFree(pReq);
804
    goto OVER;
D
dapan1121 已提交
805
  }
806 807 808 809

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

813 814
  sendInfo->msgInfo = (SDataBuf){
      .pData = pReq,
D
dapan1121 已提交
815
      .len = tlen,
816 817 818 819 820 821 822
      .handle = NULL,
  };

  sendInfo->requestId = generateRequestId();
  sendInfo->requestObjRefId = 0;
  sendInfo->param = NULL;
  sendInfo->fp = tmqHbCb;
L
Liu Jicong 已提交
823
  sendInfo->msgType = TDMT_MND_TMQ_HB;
824 825 826 827 828 829 830

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

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

OVER:
831
  taosTmrReset(tmqSendHbReq, 1000, param, tmqMgmt.timer, &tmq->hbLiveTimer);
832
  taosReleaseRef(tmqMgmt.rsetId, refId);
833 834
}

835
int32_t tmqHandleAllDelayedTask(tmq_t* pTmq) {
L
Liu Jicong 已提交
836
  STaosQall* qall = taosAllocateQall();
837
  taosReadAllQitems(pTmq->delayedTask, qall);
L
Liu Jicong 已提交
838

839 840 841 842
  if (qall->numOfItems == 0) {
    taosFreeQall(qall);
    return TSDB_CODE_SUCCESS;
  }
843

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

848
  while (pTaskType != NULL) {
849
    if (*pTaskType == TMQ_DELAYED_TASK__ASK_EP) {
850
      tmqAskEp(pTmq, true);
851 852

      int64_t* pRefId = taosMemoryMalloc(sizeof(int64_t));
853
      *pRefId = pTmq->refId;
854

X
Xiaoyu Wang 已提交
855
      tscDebug("consumer:0x%" PRIx64 " retrieve ep from mnode in 1s", pTmq->consumerId);
856
      taosTmrReset(tmqAssignAskEpTask, 1000, pRefId, tmqMgmt.timer, &pTmq->epTimer);
L
Liu Jicong 已提交
857
    } else if (*pTaskType == TMQ_DELAYED_TASK__COMMIT) {
858
      tmqCommitInner(pTmq, NULL, 1, 1, pTmq->commitCb, pTmq->commitCbUserParam);
859 860

      int64_t* pRefId = taosMemoryMalloc(sizeof(int64_t));
861
      *pRefId = pTmq->refId;
862

X
Xiaoyu Wang 已提交
863
      tscDebug("consumer:0x%" PRIx64 " commit to vnode(s) in %.2fs", pTmq->consumerId,
X
Xiaoyu Wang 已提交
864
               pTmq->autoCommitInterval / 1000.0);
865
      taosTmrReset(tmqAssignDelayedCommitTask, pTmq->autoCommitInterval, pRefId, tmqMgmt.timer, &pTmq->commitTimer);
L
Liu Jicong 已提交
866 867
    } else if (*pTaskType == TMQ_DELAYED_TASK__REPORT) {
    }
868

L
Liu Jicong 已提交
869
    taosFreeQitem(pTaskType);
870
    taosGetQitem(qall, (void**)&pTaskType);
L
Liu Jicong 已提交
871
  }
872

L
Liu Jicong 已提交
873 874 875 876
  taosFreeQall(qall);
  return 0;
}

877
static void* tmqFreeRspWrapper(SMqRspWrapper* rspWrapper) {
L
Liu Jicong 已提交
878 879 880 881 882 883 884
  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;
885 886
    taosMemoryFreeClear(pRsp->pEpset);

L
Liu Jicong 已提交
887 888 889 890 891 892
    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;
893 894
    taosMemoryFreeClear(pRsp->pEpset);

L
Liu Jicong 已提交
895 896 897
    taosMemoryFree(pRsp->metaRsp.metaRsp);
  } else if (rspWrapper->tmqRspType == TMQ_MSG_TYPE__TAOSX_RSP) {
    SMqPollRspWrapper* pRsp = (SMqPollRspWrapper*)rspWrapper;
898 899
    taosMemoryFreeClear(pRsp->pEpset);

L
Liu Jicong 已提交
900 901 902 903 904 905 906 907
    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);
  }
908 909

  return NULL;
L
Liu Jicong 已提交
910 911
}

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

L
Liu Jicong 已提交
924
  rspWrapper = NULL;
L
Liu Jicong 已提交
925 926
  taosReadAllQitems(tmq->mqueue, tmq->qall);
  while (1) {
L
Liu Jicong 已提交
927 928 929 930 931
    taosGetQitem(tmq->qall, (void**)&rspWrapper);
    if (rspWrapper) {
      tmqFreeRspWrapper(rspWrapper);
      taosFreeQitem(rspWrapper);
    } else {
L
Liu Jicong 已提交
932
      break;
L
Liu Jicong 已提交
933
    }
L
Liu Jicong 已提交
934 935 936
  }
}

D
dapan1121 已提交
937
int32_t tmqSubscribeCb(void* param, SDataBuf* pMsg, int32_t code) {
L
Liu Jicong 已提交
938 939
  SMqSubscribeCbParam* pParam = (SMqSubscribeCbParam*)param;
  pParam->rspErr = code;
dengyihao's avatar
dengyihao 已提交
940 941

  taosMemoryFree(pMsg->pEpSet);
L
Liu Jicong 已提交
942 943 944
  tsem_post(&pParam->rspSem);
  return 0;
}
945

L
Liu Jicong 已提交
946
int32_t tmq_subscription(tmq_t* tmq, tmq_list_t** topics) {
X
Xiaoyu Wang 已提交
947 948 949 950
  if (*topics == NULL) {
    *topics = tmq_list_new();
  }
  for (int i = 0; i < taosArrayGetSize(tmq->clientTopics); i++) {
L
Liu Jicong 已提交
951
    SMqClientTopic* topic = taosArrayGet(tmq->clientTopics, i);
L
Liu Jicong 已提交
952
    tmq_list_append(*topics, strchr(topic->topicName, '.') + 1);
X
Xiaoyu Wang 已提交
953
  }
L
Liu Jicong 已提交
954
  return 0;
X
Xiaoyu Wang 已提交
955 956
}

L
Liu Jicong 已提交
957
int32_t tmq_unsubscribe(tmq_t* tmq) {
L
Liu Jicong 已提交
958 959
  int32_t     rsp;
  int32_t     retryCnt = 0;
L
Liu Jicong 已提交
960
  tmq_list_t* lst = tmq_list_new();
L
Liu Jicong 已提交
961 962 963 964 965 966 967 968 969 970
  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 已提交
971 972
  tmq_list_destroy(lst);
  return rsp;
X
Xiaoyu Wang 已提交
973 974
}

975 976 977 978 979 980
static void freeClientVgImpl(void* param) {
  SMqClientTopic* pTopic = param;
  taosMemoryFreeClear(pTopic->schema.pSchema);
  taosArrayDestroy(pTopic->vgs);
}

981
void tmqFreeImpl(void* handle) {
982 983
  tmq_t*  tmq = (tmq_t*)handle;
  int64_t id = tmq->consumerId;
L
Liu Jicong 已提交
984

985
  // TODO stop timer
L
Liu Jicong 已提交
986 987 988 989
  if (tmq->mqueue) {
    tmqClearUnhandleMsg(tmq);
    taosCloseQueue(tmq->mqueue);
  }
L
Liu Jicong 已提交
990

H
Haojun Liao 已提交
991 992 993 994 995
  if (tmq->delayedTask) {
    taosCloseQueue(tmq->delayedTask);
  }

  taosFreeQall(tmq->qall);
996
  tsem_destroy(&tmq->rspSem);
H
Haojun Liao 已提交
997
  taosThreadMutexDestroy(&tmq->lock);
L
Liu Jicong 已提交
998

999
  taosArrayDestroyEx(tmq->clientTopics, freeClientVgImpl);
1000 1001
  taos_close_internal(tmq->pTscObj);
  taosMemoryFree(tmq);
1002 1003

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

1006 1007 1008 1009 1010 1011 1012 1013 1014
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);
1015
  if (tmqMgmt.rsetId < 0) {
1016 1017 1018 1019
    tmqInitRes = terrno;
  }
}

L
Liu Jicong 已提交
1020
tmq_t* tmq_consumer_new(tmq_conf_t* conf, char* errstr, int32_t errstrLen) {
1021 1022 1023 1024
  taosThreadOnce(&tmqInit, tmqMgmtInit);
  if (tmqInitRes != 0) {
    terrno = tmqInitRes;
    return NULL;
L
Liu Jicong 已提交
1025 1026
  }

L
Liu Jicong 已提交
1027 1028
  tmq_t* pTmq = taosMemoryCalloc(1, sizeof(tmq_t));
  if (pTmq == NULL) {
L
Liu Jicong 已提交
1029
    terrno = TSDB_CODE_OUT_OF_MEMORY;
1030
    tscError("failed to create consumer, groupId:%s, code:%s", conf->groupId, terrstr());
L
Liu Jicong 已提交
1031 1032
    return NULL;
  }
L
Liu Jicong 已提交
1033

L
Liu Jicong 已提交
1034 1035 1036
  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 已提交
1037 1038 1039
  pTmq->clientTopics = taosArrayInit(0, sizeof(SMqClientTopic));
  pTmq->mqueue = taosOpenQueue();
  pTmq->delayedTask = taosOpenQueue();
H
Haojun Liao 已提交
1040
  pTmq->qall = taosAllocateQall();
L
Liu Jicong 已提交
1041

H
Haojun Liao 已提交
1042
  taosThreadMutexInit(&pTmq->lock, NULL);
X
Xiaoyu Wang 已提交
1043 1044
  if (pTmq->clientTopics == NULL || pTmq->mqueue == NULL || pTmq->qall == NULL || pTmq->delayedTask == NULL ||
      conf->groupId[0] == 0) {
L
Liu Jicong 已提交
1045
    terrno = TSDB_CODE_OUT_OF_MEMORY;
1046
    tscError("consumer:0x%" PRIx64 " setup failed since %s, groupId:%s", pTmq->consumerId, terrstr(),
S
Shengliang Guan 已提交
1047
             pTmq->groupId);
1048
    goto _failed;
L
Liu Jicong 已提交
1049
  }
L
Liu Jicong 已提交
1050

L
Liu Jicong 已提交
1051 1052
  // init status
  pTmq->status = TMQ_CONSUMER_STATUS__INIT;
L
Liu Jicong 已提交
1053 1054
  pTmq->pollCnt = 0;
  pTmq->epoch = 0;
L
Liu Jicong 已提交
1055 1056
  /*pTmq->epStatus = 0;*/
  /*pTmq->epSkipCnt = 0;*/
L
Liu Jicong 已提交
1057

L
Liu Jicong 已提交
1058 1059 1060
  // set conf
  strcpy(pTmq->clientId, conf->clientId);
  strcpy(pTmq->groupId, conf->groupId);
1061
  pTmq->withTbName = conf->withTbName;
L
Liu Jicong 已提交
1062
  pTmq->useSnapshot = conf->snapEnable;
L
Liu Jicong 已提交
1063
  pTmq->autoCommit = conf->autoCommit;
L
Liu Jicong 已提交
1064
  pTmq->autoCommitInterval = conf->autoCommitInterval;
L
Liu Jicong 已提交
1065 1066
  pTmq->commitCb = conf->commitCb;
  pTmq->commitCbUserParam = conf->commitCbUserParam;
L
Liu Jicong 已提交
1067 1068
  pTmq->resetOffsetCfg = conf->resetOffset;

1069 1070
  pTmq->hbBgEnable = conf->hbBgEnable;

L
Liu Jicong 已提交
1071
  // assign consumerId
L
Liu Jicong 已提交
1072
  pTmq->consumerId = tGenIdPI64();
X
Xiaoyu Wang 已提交
1073

L
Liu Jicong 已提交
1074 1075
  // init semaphore
  if (tsem_init(&pTmq->rspSem, 0, 0) != 0) {
1076
    tscError("consumer:0x %" PRIx64 " setup failed since %s, consumer group %s", pTmq->consumerId, terrstr(),
S
Shengliang Guan 已提交
1077
             pTmq->groupId);
1078
    goto _failed;
L
Liu Jicong 已提交
1079
  }
L
Liu Jicong 已提交
1080

L
Liu Jicong 已提交
1081 1082 1083
  // init connection
  pTmq->pTscObj = taos_connect_internal(conf->ip, user, pass, NULL, NULL, conf->port, CONN_TYPE__TMQ);
  if (pTmq->pTscObj == NULL) {
1084
    tscError("consumer:0x%" PRIx64 " setup failed since %s, groupId:%s", pTmq->consumerId, terrstr(), pTmq->groupId);
L
Liu Jicong 已提交
1085
    tsem_destroy(&pTmq->rspSem);
1086
    goto _failed;
L
Liu Jicong 已提交
1087
  }
L
Liu Jicong 已提交
1088

1089 1090
  pTmq->refId = taosAddRef(tmqMgmt.rsetId, pTmq);
  if (pTmq->refId < 0) {
1091
    goto _failed;
1092 1093
  }

1094
  if (pTmq->hbBgEnable) {
L
Liu Jicong 已提交
1095 1096
    int64_t* pRefId = taosMemoryMalloc(sizeof(int64_t));
    *pRefId = pTmq->refId;
1097
    pTmq->hbLiveTimer = taosTmrStart(tmqSendHbReq, 1000, pRefId, tmqMgmt.timer);
1098 1099
  }

1100 1101 1102
  char         buf[80] = {0};
  STqOffsetVal offset = {.type = pTmq->resetOffsetCfg};
  tFormatOffset(buf, tListLen(buf), &offset);
1103 1104
  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,
1105
          pTmq->hbBgEnable);
L
Liu Jicong 已提交
1106

1107
  return pTmq;
1108

1109 1110
_failed:
  tmqFreeImpl(pTmq);
L
Liu Jicong 已提交
1111
  return NULL;
1112 1113
}

L
Liu Jicong 已提交
1114
int32_t tmq_subscribe(tmq_t* tmq, const tmq_list_t* topic_list) {
L
Liu Jicong 已提交
1115 1116 1117
  const SArray*   container = &topic_list->container;
  int32_t         sz = taosArrayGetSize(container);
  void*           buf = NULL;
L
Liu Jicong 已提交
1118
  SMsgSendInfo*   sendInfo = NULL;
L
Liu Jicong 已提交
1119
  SCMSubscribeReq req = {0};
1120
  int32_t         code = 0;
1121

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

1124
  req.consumerId = tmq->consumerId;
L
Liu Jicong 已提交
1125
  tstrncpy(req.clientId, tmq->clientId, 256);
L
Liu Jicong 已提交
1126
  tstrncpy(req.cgroup, tmq->groupId, TSDB_CGROUP_LEN);
1127 1128
  req.topicNames = taosArrayInit(sz, sizeof(void*));

1129 1130 1131 1132
  if (req.topicNames == NULL) {
    code = TSDB_CODE_OUT_OF_MEMORY;
    goto FAIL;
  }
1133

L
Liu Jicong 已提交
1134 1135
  for (int32_t i = 0; i < sz; i++) {
    char* topic = taosArrayGetP(container, i);
1136 1137

    SName name = {0};
L
Liu Jicong 已提交
1138 1139 1140 1141
    tNameSetDbName(&name, tmq->pTscObj->acctId, topic, strlen(topic));
    char* topicFName = taosMemoryCalloc(1, TSDB_TOPIC_FNAME_LEN);
    if (topicFName == NULL) {
      goto FAIL;
1142 1143
    }

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

    taosArrayPush(req.topicNames, &topicFName);
1148 1149
  }

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

L
Liu Jicong 已提交
1152
  buf = taosMemoryMalloc(tlen);
1153 1154 1155 1156
  if (buf == NULL) {
    code = TSDB_CODE_OUT_OF_MEMORY;
    goto FAIL;
  }
L
Liu Jicong 已提交
1157

1158 1159 1160
  void* abuf = buf;
  tSerializeSCMSubscribeReq(&abuf, &req);

L
Liu Jicong 已提交
1161
  sendInfo = taosMemoryCalloc(1, sizeof(SMsgSendInfo));
1162 1163 1164 1165
  if (sendInfo == NULL) {
    code = TSDB_CODE_OUT_OF_MEMORY;
    goto FAIL;
  }
1166

X
Xiaoyu Wang 已提交
1167
  SMqSubscribeCbParam param = {
L
Liu Jicong 已提交
1168
      .rspErr = 0,
1169 1170
      .refId = tmq->refId,
      .epoch = tmq->epoch,
X
Xiaoyu Wang 已提交
1171
  };
L
Liu Jicong 已提交
1172

1173 1174 1175
  if (tsem_init(&param.rspSem, 0, 0) != 0) {
    goto FAIL;
  }
L
Liu Jicong 已提交
1176 1177

  sendInfo->msgInfo = (SDataBuf){
X
Xiaoyu Wang 已提交
1178 1179 1180 1181
      .pData = buf,
      .len = tlen,
      .handle = NULL,
  };
1182

L
Liu Jicong 已提交
1183 1184
  sendInfo->requestId = generateRequestId();
  sendInfo->requestObjRefId = 0;
L
Liu Jicong 已提交
1185 1186
  sendInfo->param = &param;
  sendInfo->fp = tmqSubscribeCb;
L
Liu Jicong 已提交
1187
  sendInfo->msgType = TDMT_MND_TMQ_SUBSCRIBE;
L
Liu Jicong 已提交
1188

1189 1190 1191 1192 1193
  SEpSet epSet = getEpSet_s(&tmq->pTscObj->pAppInfo->mgmtEp);

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

L
Liu Jicong 已提交
1194 1195
  // avoid double free if msg is sent
  buf = NULL;
L
Liu Jicong 已提交
1196
  sendInfo = NULL;
L
Liu Jicong 已提交
1197

L
Liu Jicong 已提交
1198 1199
  tsem_wait(&param.rspSem);
  tsem_destroy(&param.rspSem);
1200

1201 1202 1203 1204
  if (param.rspErr != 0) {
    code = param.rspErr;
    goto FAIL;
  }
L
Liu Jicong 已提交
1205

L
Liu Jicong 已提交
1206
  int32_t retryCnt = 0;
L
Liu Jicong 已提交
1207
  while (TSDB_CODE_MND_CONSUMER_NOT_READY == tmqAskEp(tmq, false)) {
wmmhello's avatar
wmmhello 已提交
1208
    if (retryCnt++ > 40) {
L
Liu Jicong 已提交
1209 1210
      goto FAIL;
    }
1211

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

1216 1217
  // init ep timer
  if (tmq->epTimer == NULL) {
1218 1219 1220
    int64_t* pRefId1 = taosMemoryMalloc(sizeof(int64_t));
    *pRefId1 = tmq->refId;
    tmq->epTimer = taosTmrStart(tmqAssignAskEpTask, 1000, pRefId1, tmqMgmt.timer);
1221
  }
L
Liu Jicong 已提交
1222 1223

  // init auto commit timer
1224
  if (tmq->autoCommit && tmq->commitTimer == NULL) {
1225 1226 1227
    int64_t* pRefId2 = taosMemoryMalloc(sizeof(int64_t));
    *pRefId2 = tmq->refId;
    tmq->commitTimer = taosTmrStart(tmqAssignDelayedCommitTask, tmq->autoCommitInterval, pRefId2, tmqMgmt.timer);
L
Liu Jicong 已提交
1228 1229
  }

L
Liu Jicong 已提交
1230
FAIL:
L
Liu Jicong 已提交
1231
  taosArrayDestroyP(req.topicNames, taosMemoryFree);
L
Liu Jicong 已提交
1232
  taosMemoryFree(buf);
L
Liu Jicong 已提交
1233
  taosMemoryFree(sendInfo);
L
Liu Jicong 已提交
1234

L
Liu Jicong 已提交
1235
  return code;
1236 1237
}

L
Liu Jicong 已提交
1238
void tmq_conf_set_auto_commit_cb(tmq_conf_t* conf, tmq_commit_cb* cb, void* param) {
1239
  conf->commitCb = cb;
L
Liu Jicong 已提交
1240
  conf->commitCbUserParam = param;
L
Liu Jicong 已提交
1241
}
1242

D
dapan1121 已提交
1243
int32_t tmqPollCb(void* param, SDataBuf* pMsg, int32_t code) {
X
Xiaoyu Wang 已提交
1244
  SMqPollCbParam* pParam = (SMqPollCbParam*)param;
1245 1246

  int64_t         refId = pParam->refId;
X
Xiaoyu Wang 已提交
1247
  SMqClientVg*    pVg = pParam->pVg;
L
Liu Jicong 已提交
1248
  SMqClientTopic* pTopic = pParam->pTopic;
1249

1250
  tmq_t* tmq = taosAcquireRef(tmqMgmt.rsetId, refId);
1251 1252 1253 1254
  if (tmq == NULL) {
    tsem_destroy(&pParam->rspSem);
    taosMemoryFree(pParam);
    taosMemoryFree(pMsg->pData);
dengyihao's avatar
dengyihao 已提交
1255
    taosMemoryFree(pMsg->pEpSet);
1256 1257 1258 1259
    terrno = TSDB_CODE_TMQ_CONSUMER_CLOSED;
    return -1;
  }

H
Haojun Liao 已提交
1260 1261 1262 1263
  int32_t  epoch = pParam->epoch;
  int32_t  vgId = pParam->vgId;
  uint64_t requestId = pParam->requestId;

L
Liu Jicong 已提交
1264
  taosMemoryFree(pParam);
H
Haojun Liao 已提交
1265

L
Liu Jicong 已提交
1266
  if (code != 0) {
H
Haojun Liao 已提交
1267 1268
    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 已提交
1269

L
Liu Jicong 已提交
1270
    if (pMsg->pData) taosMemoryFree(pMsg->pData);
dengyihao's avatar
dengyihao 已提交
1271 1272
    if (pMsg->pEpSet) taosMemoryFree(pMsg->pEpSet);

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

L
Liu Jicong 已提交
1286 1287 1288
      pRspWrapper->tmqRspType = TMQ_MSG_TYPE__END_RSP;
      taosWriteQitem(tmq->mqueue, pRspWrapper);
    }
H
Haojun Liao 已提交
1289

L
fix txn  
Liu Jicong 已提交
1290
    goto CREATE_MSG_FAIL;
L
Liu Jicong 已提交
1291 1292
  }

X
Xiaoyu Wang 已提交
1293 1294 1295
  int32_t msgEpoch = ((SMqRspHead*)pMsg->pData)->epoch;
  int32_t tmqEpoch = atomic_load_32(&tmq->epoch);
  if (msgEpoch < tmqEpoch) {
L
Liu Jicong 已提交
1296
    // do not write into queue since updating epoch reset
H
Haojun Liao 已提交
1297 1298 1299
    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);

1300
    tsem_post(&tmq->rspSem);
1301 1302
    taosReleaseRef(tmqMgmt.rsetId, refId);

L
Liu Jicong 已提交
1303
    taosMemoryFree(pMsg->pData);
dengyihao's avatar
dengyihao 已提交
1304
    taosMemoryFree(pMsg->pEpSet);
X
Xiaoyu Wang 已提交
1305 1306 1307 1308
    return 0;
  }

  if (msgEpoch != tmqEpoch) {
H
Haojun Liao 已提交
1309 1310
    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 已提交
1311 1312
  }

L
Liu Jicong 已提交
1313 1314 1315
  // handle meta rsp
  int8_t rspType = ((SMqRspHead*)pMsg->pData)->mqMsgType;

S
Shengliang Guan 已提交
1316
  SMqPollRspWrapper* pRspWrapper = taosAllocateQitem(sizeof(SMqPollRspWrapper), DEF_QITEM, 0);
L
Liu Jicong 已提交
1317
  if (pRspWrapper == NULL) {
L
Liu Jicong 已提交
1318
    taosMemoryFree(pMsg->pData);
dengyihao's avatar
dengyihao 已提交
1319
    taosMemoryFree(pMsg->pEpSet);
H
Haojun Liao 已提交
1320
    tscWarn("consumer:0x%"PRIx64" msg discard from vgId:%d, epoch %d since out of memory", tmq->consumerId, vgId, epoch);
L
fix txn  
Liu Jicong 已提交
1321
    goto CREATE_MSG_FAIL;
L
Liu Jicong 已提交
1322
  }
L
Liu Jicong 已提交
1323

L
Liu Jicong 已提交
1324
  pRspWrapper->tmqRspType = rspType;
L
Liu Jicong 已提交
1325 1326
  pRspWrapper->vgHandle = pVg;
  pRspWrapper->topicHandle = pTopic;
H
Haojun Liao 已提交
1327
  pRspWrapper->reqId = requestId;
1328
  pRspWrapper->pEpset = pMsg->pEpSet;
L
Liu Jicong 已提交
1329

1330
  pMsg->pEpSet = NULL;
L
Liu Jicong 已提交
1331
  if (rspType == TMQ_MSG_TYPE__POLL_RSP) {
L
Liu Jicong 已提交
1332 1333 1334
    SDecoder decoder;
    tDecoderInit(&decoder, POINTER_SHIFT(pMsg->pData, sizeof(SMqRspHead)), pMsg->len - sizeof(SMqRspHead));
    tDecodeSMqDataRsp(&decoder, &pRspWrapper->dataRsp);
wmmhello's avatar
wmmhello 已提交
1335
    tDecoderClear(&decoder);
L
Liu Jicong 已提交
1336
    memcpy(&pRspWrapper->dataRsp, pMsg->pData, sizeof(SMqRspHead));
L
Liu Jicong 已提交
1337

H
Haojun Liao 已提交
1338 1339 1340 1341
    char buf[80];
    tFormatOffset(buf, 80, &pRspWrapper->dataRsp.rspOffset);
    tscDebug("consumer:0x%" PRIx64 " recv poll rsp, vgId:%d, req:%" PRId64 ", rsp:%s type %d, reqId:0x%" PRIx64,
             tmq->consumerId, vgId, pRspWrapper->dataRsp.reqOffset.version, buf, rspType, requestId);
L
Liu Jicong 已提交
1342
  } else if (rspType == TMQ_MSG_TYPE__POLL_META_RSP) {
1343 1344 1345 1346
    SDecoder decoder;
    tDecoderInit(&decoder, POINTER_SHIFT(pMsg->pData, sizeof(SMqRspHead)), pMsg->len - sizeof(SMqRspHead));
    tDecodeSMqMetaRsp(&decoder, &pRspWrapper->metaRsp);
    tDecoderClear(&decoder);
L
Liu Jicong 已提交
1347
    memcpy(&pRspWrapper->metaRsp, pMsg->pData, sizeof(SMqRspHead));
L
Liu Jicong 已提交
1348 1349 1350 1351 1352 1353
  } 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 已提交
1354 1355
  } else { // invalid rspType
    tscError("consumer:0x%"PRIx64" invalid rsp msg received, type:%d ignored", tmq->consumerId, rspType);
L
Liu Jicong 已提交
1356
  }
L
Liu Jicong 已提交
1357

L
Liu Jicong 已提交
1358
  taosMemoryFree(pMsg->pData);
H
Haojun Liao 已提交
1359
  taosWriteQitem(tmq->mqueue, pRspWrapper);
L
Liu Jicong 已提交
1360

H
Haojun Liao 已提交
1361 1362 1363
  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);

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

L
Liu Jicong 已提交
1367
  return 0;
H
Haojun Liao 已提交
1368

L
fix txn  
Liu Jicong 已提交
1369
CREATE_MSG_FAIL:
L
Liu Jicong 已提交
1370
  if (epoch == tmq->epoch) {
L
Liu Jicong 已提交
1371 1372
    atomic_store_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE);
  }
H
Haojun Liao 已提交
1373

1374
  tsem_post(&tmq->rspSem);
1375 1376
  taosReleaseRef(tmqMgmt.rsetId, refId);

L
Liu Jicong 已提交
1377
  return -1;
1378 1379
}

H
Haojun Liao 已提交
1380 1381 1382 1383 1384 1385 1386 1387 1388 1389 1390 1391 1392 1393 1394 1395 1396
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 已提交
1397 1398

    makeTopicVgroupKey(vgKey, pTopic->topicName, pVgEp->vgId);
H
Haojun Liao 已提交
1399
    STqOffsetVal* pOffset = taosHashGet(pVgOffsetHashMap, vgKey, strlen(vgKey));
H
Haojun Liao 已提交
1400

H
Haojun Liao 已提交
1401 1402 1403 1404 1405 1406 1407 1408 1409 1410 1411 1412
    STqOffsetVal  offsetNew = {.type = tmq->resetOffsetCfg};
    if (pOffset != NULL) {
      offsetNew = *pOffset;
    }

    SMqClientVg clientVg = {
        .pollCnt = 0,
        .currentOffset = offsetNew,
        .vgId = pVgEp->vgId,
        .epSet = pVgEp->epSet,
        .vgStatus = TMQ_VG_STATUS__IDLE,
        .vgSkipCnt = 0,
H
Haojun Liao 已提交
1413
        .emptyBlockReceiveTs = 0,
1414
        .numOfRows = 0,
H
Haojun Liao 已提交
1415 1416 1417 1418 1419 1420 1421 1422 1423 1424 1425 1426 1427 1428 1429 1430
    };

    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) {
1431 1432
  bool set = false;

1433
  int32_t topicNumCur = taosArrayGetSize(tmq->clientTopics);
1434
  int32_t topicNumGet = taosArrayGetSize(pRsp->topics);
1435

X
Xiaoyu Wang 已提交
1436 1437
  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",
1438
           tmq->consumerId, tmq->epoch, epoch, topicNumGet, topicNumCur);
1439 1440 1441 1442 1443 1444

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

H
Haojun Liao 已提交
1445 1446
  SHashObj* pVgOffsetHashMap = taosHashInit(64, MurmurHash3_32, false, HASH_NO_LOCK);
  if (pVgOffsetHashMap == NULL) {
1447 1448 1449
    taosArrayDestroy(newTopics);
    return false;
  }
1450

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

L
Liu Jicong 已提交
1462
        char buf[80];
L
Liu Jicong 已提交
1463
        tFormatOffset(buf, 80, &pVgCur->currentOffset);
H
Haojun Liao 已提交
1464
        tscDebug("consumer:0x%" PRIx64 ", epoch:%d vgId:%d vgKey:%s, offset:%s", tmq->consumerId, epoch,
L
Liu Jicong 已提交
1465
                 pVgCur->vgId, vgKey, buf);
H
Haojun Liao 已提交
1466
        taosHashPut(pVgOffsetHashMap, vgKey, strlen(vgKey), &pVgCur->currentOffset, sizeof(STqOffsetVal));
1467 1468 1469 1470 1471 1472 1473
      }
    }
  }

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

H
Haojun Liao 已提交
1478 1479 1480
  taosHashCleanup(pVgOffsetHashMap);

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

H
Haojun Liao 已提交
1486 1487
  tmq->clientTopics = newTopics;
  taosThreadMutexUnlock(&tmq->lock);
1488

H
Haojun Liao 已提交
1489 1490
  int8_t flag = (topicNumGet == 0)? TMQ_CONSUMER_STATUS__NO_TOPIC:TMQ_CONSUMER_STATUS__READY;
  atomic_store_8(&tmq->status, flag);
X
Xiaoyu Wang 已提交
1491
  atomic_store_32(&tmq->epoch, epoch);
H
Haojun Liao 已提交
1492

1493
  tscDebug("consumer:0x%" PRIx64 " update topic info completed", tmq->consumerId);
X
Xiaoyu Wang 已提交
1494 1495 1496
  return set;
}

H
Haojun Liao 已提交
1497
static int32_t tmqAskEpCb(void* param, SDataBuf* pMsg, int32_t code) {
1498
  SMqAskEpCbParam* pParam = (SMqAskEpCbParam*)param;
L
Liu Jicong 已提交
1499
  int8_t           async = pParam->async;
1500 1501 1502 1503 1504 1505 1506 1507 1508
  tmq_t*           tmq = taosAcquireRef(tmqMgmt.rsetId, pParam->refId);

  if (tmq == NULL) {
    if (!async) {
      tsem_destroy(&pParam->rspSem);
    } else {
      taosMemoryFree(pParam);
    }
    taosMemoryFree(pMsg->pData);
dengyihao's avatar
dengyihao 已提交
1509
    taosMemoryFree(pMsg->pEpSet);
1510 1511 1512 1513
    terrno = TSDB_CODE_TMQ_CONSUMER_CLOSED;
    return -1;
  }

L
Liu Jicong 已提交
1514
  pParam->code = code;
H
Haojun Liao 已提交
1515
  if (code != TSDB_CODE_SUCCESS) {
X
Xiaoyu Wang 已提交
1516 1517
    tscError("consumer:0x%" PRIx64 ", get topic endpoint error, async:%d, code:%s", tmq->consumerId, pParam->async,
             tstrerror(code));
L
Liu Jicong 已提交
1518
    goto END;
1519
  }
L
Liu Jicong 已提交
1520

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

L
Liu Jicong 已提交
1537
    goto END;
1538
  }
L
Liu Jicong 已提交
1539

1540 1541 1542
  tscDebug("consumer:0x%" PRIx64 ", recv ep, msg epoch %d, current epoch %d, update local ep", tmq->consumerId,
           head->epoch, epoch);

L
Liu Jicong 已提交
1543
  if (!async) {
L
Liu Jicong 已提交
1544 1545
    SMqAskEpRsp rsp;
    tDecodeSMqAskEpRsp(POINTER_SHIFT(pMsg->pData, sizeof(SMqRspHead)), &rsp);
L
Liu Jicong 已提交
1546
    tmqUpdateEp(tmq, head->epoch, &rsp);
L
Liu Jicong 已提交
1547
    tDeleteSMqAskEpRsp(&rsp);
X
Xiaoyu Wang 已提交
1548
  } else {
S
Shengliang Guan 已提交
1549
    SMqAskEpRspWrapper* pWrapper = taosAllocateQitem(sizeof(SMqAskEpRspWrapper), DEF_QITEM, 0);
L
Liu Jicong 已提交
1550
    if (pWrapper == NULL) {
X
Xiaoyu Wang 已提交
1551
      terrno = TSDB_CODE_OUT_OF_MEMORY;
L
Liu Jicong 已提交
1552 1553
      code = -1;
      goto END;
X
Xiaoyu Wang 已提交
1554
    }
1555

L
Liu Jicong 已提交
1556 1557 1558
    pWrapper->tmqRspType = TMQ_MSG_TYPE__EP_RSP;
    pWrapper->epoch = head->epoch;
    memcpy(&pWrapper->msg, pMsg->pData, sizeof(SMqRspHead));
L
Liu Jicong 已提交
1559
    tDecodeSMqAskEpRsp(POINTER_SHIFT(pMsg->pData, sizeof(SMqRspHead)), &pWrapper->msg);
L
Liu Jicong 已提交
1560

L
Liu Jicong 已提交
1561
    taosWriteQitem(tmq->mqueue, pWrapper);
1562
    tsem_post(&tmq->rspSem);
1563
  }
L
Liu Jicong 已提交
1564 1565

END:
1566 1567
  taosReleaseRef(tmqMgmt.rsetId, pParam->refId);

L
Liu Jicong 已提交
1568
  if (!async) {
L
Liu Jicong 已提交
1569
    tsem_post(&pParam->rspSem);
L
Liu Jicong 已提交
1570 1571
  } else {
    taosMemoryFree(pParam);
L
Liu Jicong 已提交
1572
  }
dengyihao's avatar
dengyihao 已提交
1573 1574

  taosMemoryFree(pMsg->pEpSet);
L
Liu Jicong 已提交
1575
  taosMemoryFree(pMsg->pData);
L
Liu Jicong 已提交
1576
  return code;
1577 1578
}

L
Liu Jicong 已提交
1579
void tmqBuildConsumeReqImpl(SMqPollReq* pReq, tmq_t* tmq, int64_t timeout, SMqClientTopic* pTopic, SMqClientVg* pVg) {
L
Liu Jicong 已提交
1580 1581 1582 1583
  int32_t groupLen = strlen(tmq->groupId);
  memcpy(pReq->subKey, tmq->groupId, groupLen);
  pReq->subKey[groupLen] = TMQ_SEPARATOR;
  strcpy(pReq->subKey + groupLen + 1, pTopic->topicName);
1584

1585
  pReq->withTbName = tmq->withTbName;
L
Liu Jicong 已提交
1586
  pReq->consumerId = tmq->consumerId;
1587
  pReq->timeout = timeout;
X
Xiaoyu Wang 已提交
1588
  pReq->epoch = tmq->epoch;
L
Liu Jicong 已提交
1589
  /*pReq->currentOffset = reqOffset;*/
L
Liu Jicong 已提交
1590
  pReq->reqOffset = pVg->currentOffset;
D
dapan1121 已提交
1591
  pReq->head.vgId = pVg->vgId;
1592 1593
  pReq->useSnapshot = tmq->useSnapshot;
  pReq->reqId = generateRequestId();
1594 1595
}

L
Liu Jicong 已提交
1596 1597
SMqMetaRspObj* tmqBuildMetaRspFromWrapper(SMqPollRspWrapper* pWrapper) {
  SMqMetaRspObj* pRspObj = taosMemoryCalloc(1, sizeof(SMqMetaRspObj));
L
Liu Jicong 已提交
1598
  pRspObj->resType = RES_TYPE__TMQ_META;
L
Liu Jicong 已提交
1599 1600 1601 1602 1603 1604 1605 1606
  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;
}

1607
SMqRspObj* tmqBuildRspFromWrapper(SMqPollRspWrapper* pWrapper, SMqClientVg* pVg, int64_t* numOfRows) {
L
Liu Jicong 已提交
1608 1609
  SMqRspObj* pRspObj = taosMemoryCalloc(1, sizeof(SMqRspObj));
  pRspObj->resType = RES_TYPE__TMQ;
1610

1611
  (*numOfRows) = 0;
L
Liu Jicong 已提交
1612 1613
  tstrncpy(pRspObj->topic, pWrapper->topicHandle->topicName, TSDB_TOPIC_FNAME_LEN);
  tstrncpy(pRspObj->db, pWrapper->topicHandle->db, TSDB_DB_FNAME_LEN);
1614

L
Liu Jicong 已提交
1615
  pRspObj->vgId = pWrapper->vgHandle->vgId;
L
Liu Jicong 已提交
1616
  pRspObj->resIter = -1;
L
Liu Jicong 已提交
1617
  memcpy(&pRspObj->rsp, &pWrapper->dataRsp, sizeof(SMqDataRsp));
L
Liu Jicong 已提交
1618

L
Liu Jicong 已提交
1619 1620
  pRspObj->resInfo.totalRows = 0;
  pRspObj->resInfo.precision = TSDB_TIME_PRECISION_MILLI;
1621

L
Liu Jicong 已提交
1622
  if (!pWrapper->dataRsp.withSchema) {
L
Liu Jicong 已提交
1623 1624
    setResSchemaInfo(&pRspObj->resInfo, pWrapper->topicHandle->schema.pSchema, pWrapper->topicHandle->schema.nCols);
  }
L
Liu Jicong 已提交
1625

1626 1627 1628 1629 1630
  // 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;
1631
    (*numOfRows) += rows;
1632 1633
  }

L
Liu Jicong 已提交
1634
  return pRspObj;
X
Xiaoyu Wang 已提交
1635 1636
}

L
Liu Jicong 已提交
1637 1638
SMqTaosxRspObj* tmqBuildTaosxRspFromWrapper(SMqPollRspWrapper* pWrapper) {
  SMqTaosxRspObj* pRspObj = taosMemoryCalloc(1, sizeof(SMqTaosxRspObj));
1639
  pRspObj->resType = RES_TYPE__TMQ_METADATA;
L
Liu Jicong 已提交
1640 1641 1642 1643
  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;
1644
  memcpy(&pRspObj->rsp, &pWrapper->taosxRsp, sizeof(STaosxRsp));
L
Liu Jicong 已提交
1645 1646 1647

  pRspObj->resInfo.totalRows = 0;
  pRspObj->resInfo.precision = TSDB_TIME_PRECISION_MILLI;
1648
  if (!pWrapper->taosxRsp.withSchema) {
L
Liu Jicong 已提交
1649 1650 1651 1652 1653 1654
    setResSchemaInfo(&pRspObj->resInfo, pWrapper->topicHandle->schema.pSchema, pWrapper->topicHandle->schema.nCols);
  }

  return pRspObj;
}

1655 1656 1657 1658 1659 1660 1661 1662 1663 1664 1665 1666 1667 1668 1669 1670 1671 1672 1673 1674 1675 1676 1677 1678 1679 1680 1681 1682 1683 1684 1685 1686 1687
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 已提交
1688
  pParam->pVg = pVg;  // pVg may be released,fix it
1689 1690
  pParam->pTopic = pTopic;
  pParam->vgId = pVg->vgId;
H
Haojun Liao 已提交
1691
  pParam->requestId = req.reqId;
1692 1693 1694 1695 1696 1697 1698 1699 1700 1701 1702 1703 1704 1705 1706 1707 1708 1709 1710 1711 1712 1713 1714 1715

  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 已提交
1716
  tscDebug("consumer:0x%" PRIx64 " send poll to %s vgId:%d, epoch %d, req:%s, reqId:0x%" PRIx64,
1717 1718 1719 1720 1721 1722 1723 1724 1725
           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;
}

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

  for (int i = 0; i < numOfTopics; i++) {
X
Xiaoyu Wang 已提交
1732
    SMqClientTopic* pTopic = taosArrayGet(tmq->clientTopics, i);
X
Xiaoyu Wang 已提交
1733
    int32_t         numOfVg = taosArrayGetSize(pTopic->vgs);
1734 1735

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

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

L
Liu Jicong 已提交
1758
      atomic_store_32(&pVg->vgSkipCnt, 0);
1759 1760 1761
      int32_t code = doTmqPollImpl(tmq, pTopic, pVg, timeout);
      if (code != TSDB_CODE_SUCCESS) {
        return code;
D
dapan1121 已提交
1762
      }
X
Xiaoyu Wang 已提交
1763 1764
    }
  }
1765

1766
  tscDebug("consumer:0x%" PRIx64 " end to poll data", tmq->consumerId);
X
Xiaoyu Wang 已提交
1767 1768 1769
  return 0;
}

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

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

X
Xiaoyu Wang 已提交
1793
  while (1) {
1794 1795
    SMqRspWrapper* pRspWrapper = NULL;
    taosGetQitem(tmq->qall, (void**)&pRspWrapper);
1796

1797
    if (pRspWrapper == NULL) {
L
Liu Jicong 已提交
1798
      taosReadAllQitems(tmq->mqueue, tmq->qall);
1799
      taosGetQitem(tmq->qall, (void**)&pRspWrapper);
L
Liu Jicong 已提交
1800

1801
      if (pRspWrapper == NULL) {
L
Liu Jicong 已提交
1802 1803
        return NULL;
      }
X
Xiaoyu Wang 已提交
1804 1805
    }

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

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

1816
      int32_t consumerEpoch = atomic_load_32(&tmq->epoch);
1817 1818 1819 1820
      SMqDataRsp* pDataRsp = &pollRspWrapper->dataRsp;

      if (pDataRsp->head.epoch == consumerEpoch) {
        // todo fix it: race condition
L
Liu Jicong 已提交
1821
        SMqClientVg* pVg = pollRspWrapper->vgHandle;
1822 1823 1824 1825 1826 1827 1828 1829 1830 1831 1832

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

        pVg->currentOffset = pDataRsp->rspOffset;
X
Xiaoyu Wang 已提交
1833
        atomic_store_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE);
H
Haojun Liao 已提交
1834

1835 1836 1837 1838 1839
        char buf[80];
        tFormatOffset(buf, 80, &pDataRsp->rspOffset);
        if (pDataRsp->blockNum == 0) {
          tscDebug("consumer:0x%" PRIx64 " empty block received, vgId:%d, offset:%s, reqId:0x%" PRIx64, tmq->consumerId,
                   pVg->vgId, buf, pollRspWrapper->reqId);
1840
          pRspWrapper = tmqFreeRspWrapper(pRspWrapper);
L
Liu Jicong 已提交
1841
          taosFreeQitem(pollRspWrapper);
1842
        } else {  // build rsp
1843 1844
          int64_t numOfRows = 0;
          SMqRspObj* pRsp = tmqBuildRspFromWrapper(pollRspWrapper, pVg, &numOfRows);
H
Haojun Liao 已提交
1845 1846 1847 1848
          tscDebug("consumer:0x%" PRIx64 " process poll rsp, vgId:%d, offset:%s, blocks:%d, rows:%" PRId64
                   ", vg total:%" PRId64 " reqId:0x%" PRIx64,
                   tmq->consumerId, pVg->vgId, buf, pDataRsp->blockNum, numOfRows, pVg->numOfRows,
                   pollRspWrapper->reqId);
H
Haojun Liao 已提交
1849

1850
          tmq->totalRows += numOfRows;
1851 1852 1853
          taosFreeQitem(pollRspWrapper);
          return pRsp;
        }
X
Xiaoyu Wang 已提交
1854
      } else {
X
Xiaoyu Wang 已提交
1855
        tscDebug("consumer:0x%" PRIx64 " msg discard since epoch mismatch: msg epoch %d, consumer epoch %d",
1856
                 tmq->consumerId, pDataRsp->head.epoch, consumerEpoch);
1857
        pRspWrapper = tmqFreeRspWrapper(pRspWrapper);
L
Liu Jicong 已提交
1858 1859
        taosFreeQitem(pollRspWrapper);
      }
1860 1861
    } else if (pRspWrapper->tmqRspType == TMQ_MSG_TYPE__POLL_META_RSP) {
      SMqPollRspWrapper* pollRspWrapper = (SMqPollRspWrapper*)pRspWrapper;
L
Liu Jicong 已提交
1862
      int32_t            consumerEpoch = atomic_load_32(&tmq->epoch);
1863 1864 1865

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

L
Liu Jicong 已提交
1866
      if (pollRspWrapper->metaRsp.head.epoch == consumerEpoch) {
L
Liu Jicong 已提交
1867
        SMqClientVg* pVg = pollRspWrapper->vgHandle;
wmmhello's avatar
wmmhello 已提交
1868
        pVg->currentOffset = pollRspWrapper->metaRsp.rspOffset;
L
Liu Jicong 已提交
1869 1870
        atomic_store_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE);
        // build rsp
L
Liu Jicong 已提交
1871
        SMqMetaRspObj* pRsp = tmqBuildMetaRspFromWrapper(pollRspWrapper);
L
Liu Jicong 已提交
1872 1873 1874
        taosFreeQitem(pollRspWrapper);
        return pRsp;
      } else {
X
Xiaoyu Wang 已提交
1875
        tscDebug("consumer:0x%" PRIx64 " msg discard since epoch mismatch: msg epoch %d, consumer epoch %d",
1876
                 tmq->consumerId, pollRspWrapper->metaRsp.head.epoch, consumerEpoch);
1877
        pRspWrapper = tmqFreeRspWrapper(pRspWrapper);
L
Liu Jicong 已提交
1878
        taosFreeQitem(pollRspWrapper);
X
Xiaoyu Wang 已提交
1879
      }
1880 1881
    } else if (pRspWrapper->tmqRspType == TMQ_MSG_TYPE__TAOSX_RSP) {
      SMqPollRspWrapper* pollRspWrapper = (SMqPollRspWrapper*)pRspWrapper;
L
Liu Jicong 已提交
1882
      int32_t consumerEpoch = atomic_load_32(&tmq->epoch);
H
Haojun Liao 已提交
1883

L
Liu Jicong 已提交
1884 1885 1886 1887
      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 已提交
1888

L
Liu Jicong 已提交
1889
        if (pollRspWrapper->taosxRsp.blockNum == 0) {
H
Haojun Liao 已提交
1890 1891
          tscDebug("consumer:0x%" PRIx64 " taosx empty block received, vgId:%d, reqId:0x%" PRIx64, tmq->consumerId, pVg->vgId,
                   pollRspWrapper->reqId);
H
Haojun Liao 已提交
1892
          pVg->emptyBlockReceiveTs = taosGetTimestampMs();
1893
          pRspWrapper = tmqFreeRspWrapper(pRspWrapper);
H
Haojun Liao 已提交
1894
          taosFreeQitem(pollRspWrapper);
L
Liu Jicong 已提交
1895
          continue;
H
Haojun Liao 已提交
1896 1897
        } else {
          pVg->emptyBlockReceiveTs = 0; // reset the ts
L
Liu Jicong 已提交
1898
        }
wmmhello's avatar
wmmhello 已提交
1899

L
Liu Jicong 已提交
1900
        // build rsp
wmmhello's avatar
wmmhello 已提交
1901
        void* pRsp = NULL;
1902
        int64_t numOfRows = 0;
L
Liu Jicong 已提交
1903
        if (pollRspWrapper->taosxRsp.createTableNum == 0) {
1904
          pRsp = tmqBuildRspFromWrapper(pollRspWrapper, pVg, &numOfRows);
L
Liu Jicong 已提交
1905
        } else {
wmmhello's avatar
wmmhello 已提交
1906 1907
          pRsp = tmqBuildTaosxRspFromWrapper(pollRspWrapper);
        }
H
Haojun Liao 已提交
1908

1909 1910
        tmq->totalRows += numOfRows;

H
Haojun Liao 已提交
1911 1912
        char buf[80];
        tFormatOffset(buf, 80, &pVg->currentOffset);
H
Haojun Liao 已提交
1913 1914 1915 1916
        tscDebug("consumer:0x%" PRIx64 " process taosx poll rsp, vgId:%d, offset:%s, blocks:%d, rows:%" PRId64
                 ", vg total:%" PRId64 " reqId:0x%" PRIx64,
                 tmq->consumerId, pVg->vgId, buf, pollRspWrapper->dataRsp.blockNum, numOfRows, pVg->numOfRows,
                 pollRspWrapper->reqId);
H
Haojun Liao 已提交
1917 1918

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

L
Liu Jicong 已提交
1921
      } else {
X
Xiaoyu Wang 已提交
1922
        tscDebug("consumer:0x%" PRIx64 " msg discard since epoch mismatch: msg epoch %d, consumer epoch %d",
1923
                 tmq->consumerId, pollRspWrapper->taosxRsp.head.epoch, consumerEpoch);
1924
        pRspWrapper = tmqFreeRspWrapper(pRspWrapper);
L
Liu Jicong 已提交
1925 1926
        taosFreeQitem(pollRspWrapper);
      }
X
Xiaoyu Wang 已提交
1927
    } else {
H
Haojun Liao 已提交
1928 1929
      tscDebug("consumer:0x%" PRIx64 " not data msg received", tmq->consumerId);

X
Xiaoyu Wang 已提交
1930
      bool reset = false;
1931 1932
      tmqHandleNoPollRsp(tmq, pRspWrapper, &reset);
      taosFreeQitem(pRspWrapper);
X
Xiaoyu Wang 已提交
1933
      if (pollIfReset && reset) {
1934
        tscDebug("consumer:0x%" PRIx64 ", reset and repoll", tmq->consumerId);
1935
        tmqPollImpl(tmq, timeout);
X
Xiaoyu Wang 已提交
1936 1937 1938 1939 1940
      }
    }
  }
}

1941
TAOS_RES* tmq_consumer_poll(tmq_t* tmq, int64_t timeout) {
L
Liu Jicong 已提交
1942 1943
  void*   rspObj;
  int64_t startTime = taosGetTimestampMs();
L
Liu Jicong 已提交
1944

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

1947 1948 1949
#if 0
  tmqHandleAllDelayedTask(tmq);
  tmqPollImpl(tmq, timeout);
1950
  rspObj = tmqHandleAllRsp(tmq, timeout, false);
L
Liu Jicong 已提交
1951 1952
  if (rspObj) {
    return (TAOS_RES*)rspObj;
L
fix  
Liu Jicong 已提交
1953
  }
1954
#endif
X
Xiaoyu Wang 已提交
1955

1956
  // in no topic status, delayed task also need to be processed
L
Liu Jicong 已提交
1957
  if (atomic_load_8(&tmq->status) == TMQ_CONSUMER_STATUS__INIT) {
1958
    tscDebug("consumer:0x%" PRIx64 " poll return since consumer is init", tmq->consumerId);
1959
    taosMsleep(500);  //     sleep for a while
1960 1961 1962
    return NULL;
  }

wmmhello's avatar
wmmhello 已提交
1963
  while (atomic_load_8(&tmq->status) == TMQ_CONSUMER_STATUS__RECOVER) {
L
Liu Jicong 已提交
1964 1965
    int32_t retryCnt = 0;
    while (TSDB_CODE_MND_CONSUMER_NOT_READY == tmqAskEp(tmq, false)) {
H
Haojun Liao 已提交
1966
      if (retryCnt++ > 40) {
L
Liu Jicong 已提交
1967 1968
        return NULL;
      }
1969

H
Haojun Liao 已提交
1970
      tscDebug("consumer:0x%" PRIx64 " not ready, retry:%d/40 in 500ms", tmq->consumerId, retryCnt);
L
Liu Jicong 已提交
1971 1972 1973 1974
      taosMsleep(500);
    }
  }

X
Xiaoyu Wang 已提交
1975
  while (1) {
L
Liu Jicong 已提交
1976
    tmqHandleAllDelayedTask(tmq);
1977

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

1982
    rspObj = tmqHandleAllRsp(tmq, timeout, false);
L
Liu Jicong 已提交
1983
    if (rspObj) {
1984
      tscDebug("consumer:0x%" PRIx64 " return rsp %p", tmq->consumerId, rspObj);
L
Liu Jicong 已提交
1985
      return (TAOS_RES*)rspObj;
L
Liu Jicong 已提交
1986
    } else if (terrno == TSDB_CODE_TQ_NO_COMMITTED_OFFSET) {
1987
      tscDebug("consumer:0x%" PRIx64 " return null since no committed offset", tmq->consumerId);
L
Liu Jicong 已提交
1988
      return NULL;
X
Xiaoyu Wang 已提交
1989
    }
1990

1991
    if (timeout >= 0) {
L
Liu Jicong 已提交
1992
      int64_t currentTime = taosGetTimestampMs();
1993 1994 1995
      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 已提交
1996
                 tmq->consumerId, tmq->epoch, startTime, currentTime);
X
Xiaoyu Wang 已提交
1997 1998
        return NULL;
      }
1999
      tsem_timewait(&tmq->rspSem, (timeout - elapsedTime));
L
Liu Jicong 已提交
2000 2001
    } else {
      // use tsem_timewait instead of tsem_wait to avoid unexpected stuck
L
Liu Jicong 已提交
2002
      tsem_timewait(&tmq->rspSem, 1000);
X
Xiaoyu Wang 已提交
2003 2004 2005 2006
    }
  }
}

L
Liu Jicong 已提交
2007
int32_t tmq_consumer_close(tmq_t* tmq) {
H
Haojun Liao 已提交
2008 2009
  tscDebug("consumer:0x%" PRIx64" start to close consumer, status:%d", tmq->consumerId, tmq->status);

2010
  if (tmq->status == TMQ_CONSUMER_STATUS__READY) {
L
Liu Jicong 已提交
2011 2012
    int32_t rsp = tmq_commit_sync(tmq, NULL);
    if (rsp != 0) {
L
Liu Jicong 已提交
2013
      return rsp;
2014 2015
    }

2016 2017 2018 2019 2020 2021 2022 2023 2024 2025 2026 2027 2028 2029 2030 2031 2032 2033
    int32_t numOfTopics = taosArrayGetSize(tmq->clientTopics);
    tscDebug("consumer:0x%" PRIx64 " closing poll:%" PRId64 " rows:%" PRId64 " topics:%d, final epoch:%d",
             tmq->consumerId, tmq->pollCnt, tmq->totalRows, numOfTopics, tmq->epoch);

    tscDebug("consumer:0x%" PRIx64 " rows dist begin: ", tmq->consumerId);
    for (int32_t i = 0; i < numOfTopics; ++i) {
      SMqClientTopic* pTopics = taosArrayGet(tmq->clientTopics, i);

      tscDebug("consumer:0x%" PRIx64 " topic:%d", tmq->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);
      }
    }

    tscDebug("consumer:0x%" PRIx64 " rows dist end", tmq->consumerId);

L
Liu Jicong 已提交
2034
    int32_t     retryCnt = 0;
2035
    tmq_list_t* lst = tmq_list_new();
L
Liu Jicong 已提交
2036 2037 2038 2039 2040 2041 2042 2043 2044 2045
    while (1) {
      rsp = tmq_subscribe(tmq, lst);
      if (rsp != TSDB_CODE_MND_CONSUMER_NOT_READY || retryCnt > 5) {
        break;
      } else {
        retryCnt++;
        taosMsleep(500);
      }
    }

2046
    tmq_list_destroy(lst);
L
Liu Jicong 已提交
2047
  }
H
Haojun Liao 已提交
2048

2049
  taosRemoveRef(tmqMgmt.rsetId, tmq->refId);
L
Liu Jicong 已提交
2050
  return 0;
2051
}
L
Liu Jicong 已提交
2052

L
Liu Jicong 已提交
2053 2054
const char* tmq_err2str(int32_t err) {
  if (err == 0) {
L
Liu Jicong 已提交
2055
    return "success";
L
Liu Jicong 已提交
2056
  } else if (err == -1) {
L
Liu Jicong 已提交
2057 2058 2059
    return "fail";
  } else {
    return tstrerror(err);
L
Liu Jicong 已提交
2060 2061
  }
}
L
Liu Jicong 已提交
2062

L
Liu Jicong 已提交
2063 2064 2065 2066 2067
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;
2068 2069
  } else if (TD_RES_TMQ_METADATA(res)) {
    return TMQ_RES_METADATA;
L
Liu Jicong 已提交
2070 2071 2072 2073 2074
  } else {
    return TMQ_RES_INVALID;
  }
}

L
Liu Jicong 已提交
2075
const char* tmq_get_topic_name(TAOS_RES* res) {
L
Liu Jicong 已提交
2076 2077
  if (TD_RES_TMQ(res)) {
    SMqRspObj* pRspObj = (SMqRspObj*)res;
L
Liu Jicong 已提交
2078
    return strchr(pRspObj->topic, '.') + 1;
L
Liu Jicong 已提交
2079 2080 2081
  } else if (TD_RES_TMQ_META(res)) {
    SMqMetaRspObj* pMetaRspObj = (SMqMetaRspObj*)res;
    return strchr(pMetaRspObj->topic, '.') + 1;
2082 2083 2084
  } else if (TD_RES_TMQ_METADATA(res)) {
    SMqTaosxRspObj* pRspObj = (SMqTaosxRspObj*)res;
    return strchr(pRspObj->topic, '.') + 1;
L
Liu Jicong 已提交
2085 2086 2087 2088 2089
  } else {
    return NULL;
  }
}

L
Liu Jicong 已提交
2090 2091 2092 2093
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 已提交
2094 2095 2096
  } else if (TD_RES_TMQ_META(res)) {
    SMqMetaRspObj* pMetaRspObj = (SMqMetaRspObj*)res;
    return strchr(pMetaRspObj->db, '.') + 1;
2097 2098 2099
  } else if (TD_RES_TMQ_METADATA(res)) {
    SMqTaosxRspObj* pRspObj = (SMqTaosxRspObj*)res;
    return strchr(pRspObj->db, '.') + 1;
L
Liu Jicong 已提交
2100 2101 2102 2103 2104
  } else {
    return NULL;
  }
}

L
Liu Jicong 已提交
2105 2106 2107 2108
int32_t tmq_get_vgroup_id(TAOS_RES* res) {
  if (TD_RES_TMQ(res)) {
    SMqRspObj* pRspObj = (SMqRspObj*)res;
    return pRspObj->vgId;
L
Liu Jicong 已提交
2109 2110 2111
  } else if (TD_RES_TMQ_META(res)) {
    SMqMetaRspObj* pMetaRspObj = (SMqMetaRspObj*)res;
    return pMetaRspObj->vgId;
2112
  } else if (TD_RES_TMQ_METADATA(res)) {
2113 2114
    SMqTaosxRspObj* pRspObj = (SMqTaosxRspObj*)res;
    return pRspObj->vgId;
L
Liu Jicong 已提交
2115 2116 2117 2118
  } else {
    return -1;
  }
}
L
Liu Jicong 已提交
2119 2120 2121 2122 2123 2124 2125 2126

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;
    }
2127
    return (const char*)taosArrayGetP(pRspObj->rsp.blockTbName, pRspObj->resIter);
2128 2129
  } else if (TD_RES_TMQ_METADATA(res)) {
    SMqTaosxRspObj* pRspObj = (SMqTaosxRspObj*)res;
L
Liu Jicong 已提交
2130 2131 2132
    if (!pRspObj->rsp.withTbName || pRspObj->rsp.blockTbName == NULL || pRspObj->resIter < 0 ||
        pRspObj->resIter >= pRspObj->rsp.blockNum) {
      return NULL;
2133
    }
L
Liu Jicong 已提交
2134 2135
    return (const char*)taosArrayGetP(pRspObj->rsp.blockTbName, pRspObj->resIter);
  }
L
Liu Jicong 已提交
2136 2137
  return NULL;
}
2138

L
Liu Jicong 已提交
2139
void tmq_commit_async(tmq_t* tmq, const TAOS_RES* msg, tmq_commit_cb* cb, void* param) {
L
Liu Jicong 已提交
2140
  tmqCommitInner(tmq, msg, 0, 1, cb, param);
L
Liu Jicong 已提交
2141 2142
}

2143
int32_t tmq_commit_sync(tmq_t* tmq, const TAOS_RES* msg) {
L
Liu Jicong 已提交
2144
  return tmqCommitInner(tmq, msg, 0, 0, NULL, NULL);
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 2199 2200 2201 2202 2203 2204 2205 2206 2207 2208 2209 2210 2211 2212 2213 2214 2215 2216 2217 2218 2219 2220 2221 2222 2223 2224 2225 2226 2227 2228 2229 2230 2231 2232 2233 2234 2235

int32_t tmqAskEp(tmq_t* tmq, bool async) {
  int32_t code = TSDB_CODE_SUCCESS;
#if 0
  int8_t  epStatus = atomic_val_compare_exchange_8(&tmq->epStatus, 0, 1);
  if (epStatus == 1) {
    int32_t epSkipCnt = atomic_add_fetch_32(&tmq->epSkipCnt, 1);
    tscTrace("consumer:0x%" PRIx64 ", skip ask ep cnt %d", tmq->consumerId, epSkipCnt);
    if (epSkipCnt < 5000) return 0;
  }
  atomic_store_32(&tmq->epSkipCnt, 0);
#endif

  SMqAskEpReq req = {0};
  req.consumerId = tmq->consumerId;
  req.epoch = tmq->epoch;
  strcpy(req.cgroup, tmq->groupId);

  int32_t tlen = tSerializeSMqAskEpReq(NULL, 0, &req);
  if (tlen < 0) {
    tscError("consumer:0x%" PRIx64 ", tSerializeSMqAskEpReq failed", tmq->consumerId);
    return -1;
  }

  void* pReq = taosMemoryCalloc(1, tlen);
  if (pReq == NULL) {
    tscError("consumer:0x%" PRIx64 ", failed to malloc askEpReq msg, size:%d", tmq->consumerId, tlen);
    terrno = TSDB_CODE_OUT_OF_MEMORY;
    return -1;
  }

  if (tSerializeSMqAskEpReq(pReq, tlen, &req) < 0) {
    tscError("consumer:0x%" PRIx64 ", tSerializeSMqAskEpReq %d failed", tmq->consumerId, tlen);
    taosMemoryFree(pReq);
    return -1;
  }

  SMqAskEpCbParam* pParam = taosMemoryCalloc(1, sizeof(SMqAskEpCbParam));
  if (pParam == NULL) {
    tscError("consumer:0x%" PRIx64 ", failed to malloc subscribe param", tmq->consumerId);
    taosMemoryFree(pReq);
    return -1;
  }

  pParam->refId = tmq->refId;
  pParam->epoch = tmq->epoch;
  pParam->async = async;
  tsem_init(&pParam->rspSem, 0, 0);

  SMsgSendInfo* sendInfo = taosMemoryCalloc(1, sizeof(SMsgSendInfo));
  if (sendInfo == NULL) {
    tsem_destroy(&pParam->rspSem);
    taosMemoryFree(pParam);
    taosMemoryFree(pReq);
    return -1;
  }

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

  sendInfo->requestId = generateRequestId();
  sendInfo->requestObjRefId = 0;
  sendInfo->param = pParam;
  sendInfo->fp = tmqAskEpCb;
  sendInfo->msgType = TDMT_MND_TMQ_ASK_EP;

  SEpSet epSet = getEpSet_s(&tmq->pTscObj->pAppInfo->mgmtEp);
  tscDebug("consumer:0x%" PRIx64 " ask ep from mnode, async:%d, reqId:0x%" PRIx64, tmq->consumerId, async,
           sendInfo->requestId);

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

  if (!async) {
    tsem_wait(&pParam->rspSem);
    code = pParam->code;
    taosMemoryFree(pParam);
  }

  return code;
}

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

int32_t tmqCommitDone(SMqCommitCbParamSet* pParamSet) {
2236 2237 2238
  int64_t refId = pParamSet->refId;

  tmq_t* tmq = taosAcquireRef(tmqMgmt.rsetId, refId);
2239 2240 2241 2242 2243 2244 2245 2246 2247 2248 2249 2250 2251 2252
  if (tmq == NULL) {
    if (!pParamSet->async) {
      tsem_destroy(&pParamSet->rspSem);
    }
    taosMemoryFree(pParamSet);
    terrno = TSDB_CODE_TMQ_CONSUMER_CLOSED;
    return -1;
  }

  // if no more waiting rsp
  if (pParamSet->async) {
    // call async cb func
    if (pParamSet->automatic && tmq->commitCb) {
      tmq->commitCb(tmq, pParamSet->rspErr, tmq->commitCbUserParam);
2253
    } else if (!pParamSet->automatic && pParamSet->userCb) { // sem post
2254 2255
      pParamSet->userCb(tmq, pParamSet->rspErr, pParamSet->userParam);
    }
2256

2257 2258 2259 2260 2261 2262 2263 2264 2265
    taosMemoryFree(pParamSet);
  } else {
    tsem_post(&pParamSet->rspSem);
  }

#if 0
  taosArrayDestroyP(pParamSet->successfulOffsets, taosMemoryFree);
    taosArrayDestroyP(pParamSet->failedOffsets, taosMemoryFree);
#endif
2266 2267

  taosReleaseRef(tmqMgmt.rsetId, refId);
2268
  return 0;
2269 2270 2271 2272 2273
}

void tmqCommitRspCountDown(SMqCommitCbParamSet* pParamSet, int64_t consumerId, const char* pTopic, int32_t vgId) {
  int32_t waitingRspNum = atomic_sub_fetch_32(&pParamSet->waitingRspNum, 1);
  if (waitingRspNum == 0) {
H
Haojun Liao 已提交
2274 2275
    tscDebug("consumer:0x%" PRIx64 " topic:%s vgId:%d all commit-rsp received, commit completed", consumerId, pTopic,
             vgId);
2276
    tmqCommitDone(pParamSet);
H
Haojun Liao 已提交
2277 2278 2279
  } else {
    tscDebug("consumer:0x%" PRIx64 " topic:%s vgId:%d commit-rsp received, remain:%d", consumerId, pTopic, vgId,
             waitingRspNum);
2280 2281
  }
}