clientTmq.c 71.2 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
typedef struct SVgroupSaveInfo {
  STqOffsetVal offset;
  int64_t      numOfRows;
} SVgroupSaveInfo;

H
Haojun Liao 已提交
1385 1386 1387 1388 1389 1390 1391 1392 1393 1394 1395 1396 1397 1398 1399 1400 1401
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 已提交
1402 1403

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

H
Haojun Liao 已提交
1406
    int64_t numOfRows = 0;
H
Haojun Liao 已提交
1407
    STqOffsetVal  offsetNew = {.type = tmq->resetOffsetCfg};
H
Haojun Liao 已提交
1408 1409 1410
    if (pInfo != NULL) {
      offsetNew = pInfo->offset;
      numOfRows = pInfo->numOfRows;
H
Haojun Liao 已提交
1411 1412 1413 1414 1415 1416 1417 1418 1419
    }

    SMqClientVg clientVg = {
        .pollCnt = 0,
        .currentOffset = offsetNew,
        .vgId = pVgEp->vgId,
        .epSet = pVgEp->epSet,
        .vgStatus = TMQ_VG_STATUS__IDLE,
        .vgSkipCnt = 0,
H
Haojun Liao 已提交
1420
        .emptyBlockReceiveTs = 0,
H
Haojun Liao 已提交
1421
        .numOfRows = numOfRows,
H
Haojun Liao 已提交
1422 1423 1424 1425 1426 1427 1428 1429 1430 1431 1432 1433 1434 1435 1436 1437
    };

    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) {
1438 1439
  bool set = false;

1440
  int32_t topicNumCur = taosArrayGetSize(tmq->clientTopics);
1441
  int32_t topicNumGet = taosArrayGetSize(pRsp->topics);
1442

X
Xiaoyu Wang 已提交
1443 1444
  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",
1445
           tmq->consumerId, tmq->epoch, epoch, topicNumGet, topicNumCur);
1446 1447 1448 1449 1450 1451

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

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

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

L
Liu Jicong 已提交
1469
        char buf[80];
L
Liu Jicong 已提交
1470
        tFormatOffset(buf, 80, &pVgCur->currentOffset);
H
Haojun Liao 已提交
1471
        tscDebug("consumer:0x%" PRIx64 ", epoch:%d vgId:%d vgKey:%s, offset:%s", tmq->consumerId, epoch,
L
Liu Jicong 已提交
1472
                 pVgCur->vgId, vgKey, buf);
H
Haojun Liao 已提交
1473 1474 1475

        SVgroupSaveInfo info = {.offset = pVgCur->currentOffset, .numOfRows = pVgCur->numOfRows};
        taosHashPut(pVgOffsetHashMap, vgKey, strlen(vgKey), &info, sizeof(SVgroupSaveInfo));
1476 1477 1478 1479 1480 1481 1482
      }
    }
  }

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

H
Haojun Liao 已提交
1487 1488 1489
  taosHashCleanup(pVgOffsetHashMap);

  taosThreadMutexLock(&tmq->lock);
1490
  // destroy current buffered existed topics info
1491
  if (tmq->clientTopics) {
H
Haojun Liao 已提交
1492
    taosArrayDestroyEx(tmq->clientTopics, freeClientVgInfo);
X
Xiaoyu Wang 已提交
1493
  }
1494

H
Haojun Liao 已提交
1495 1496
  tmq->clientTopics = newTopics;
  taosThreadMutexUnlock(&tmq->lock);
1497

H
Haojun Liao 已提交
1498 1499
  int8_t flag = (topicNumGet == 0)? TMQ_CONSUMER_STATUS__NO_TOPIC:TMQ_CONSUMER_STATUS__READY;
  atomic_store_8(&tmq->status, flag);
X
Xiaoyu Wang 已提交
1500
  atomic_store_32(&tmq->epoch, epoch);
H
Haojun Liao 已提交
1501

1502
  tscDebug("consumer:0x%" PRIx64 " update topic info completed", tmq->consumerId);
X
Xiaoyu Wang 已提交
1503 1504 1505
  return set;
}

H
Haojun Liao 已提交
1506
static int32_t tmqAskEpCb(void* param, SDataBuf* pMsg, int32_t code) {
1507
  SMqAskEpCbParam* pParam = (SMqAskEpCbParam*)param;
L
Liu Jicong 已提交
1508
  int8_t           async = pParam->async;
1509 1510 1511 1512 1513 1514 1515 1516 1517
  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 已提交
1518
    taosMemoryFree(pMsg->pEpSet);
1519 1520 1521 1522
    terrno = TSDB_CODE_TMQ_CONSUMER_CLOSED;
    return -1;
  }

L
Liu Jicong 已提交
1523
  pParam->code = code;
H
Haojun Liao 已提交
1524
  if (code != TSDB_CODE_SUCCESS) {
X
Xiaoyu Wang 已提交
1525 1526
    tscError("consumer:0x%" PRIx64 ", get topic endpoint error, async:%d, code:%s", tmq->consumerId, pParam->async,
             tstrerror(code));
L
Liu Jicong 已提交
1527
    goto END;
1528
  }
L
Liu Jicong 已提交
1529

L
Liu Jicong 已提交
1530
  // tmq's epoch is monotonically increase,
L
Liu Jicong 已提交
1531
  // so it's safe to discard any old epoch msg.
L
Liu Jicong 已提交
1532
  // Epoch will only increase when received newer epoch ep msg
L
Liu Jicong 已提交
1533 1534 1535
  SMqRspHead* head = pMsg->pData;
  int32_t     epoch = atomic_load_32(&tmq->epoch);
  if (head->epoch <= epoch) {
1536 1537
    tscDebug("consumer:0x%" PRIx64 ", recv ep, msg epoch %d, current epoch %d, no need to update local ep",
             tmq->consumerId, head->epoch, epoch);
1538 1539 1540 1541 1542 1543 1544 1545
    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 已提交
1546
    goto END;
1547
  }
L
Liu Jicong 已提交
1548

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

L
Liu Jicong 已提交
1552
  if (!async) {
L
Liu Jicong 已提交
1553 1554
    SMqAskEpRsp rsp;
    tDecodeSMqAskEpRsp(POINTER_SHIFT(pMsg->pData, sizeof(SMqRspHead)), &rsp);
L
Liu Jicong 已提交
1555
    tmqUpdateEp(tmq, head->epoch, &rsp);
L
Liu Jicong 已提交
1556
    tDeleteSMqAskEpRsp(&rsp);
X
Xiaoyu Wang 已提交
1557
  } else {
S
Shengliang Guan 已提交
1558
    SMqAskEpRspWrapper* pWrapper = taosAllocateQitem(sizeof(SMqAskEpRspWrapper), DEF_QITEM, 0);
L
Liu Jicong 已提交
1559
    if (pWrapper == NULL) {
X
Xiaoyu Wang 已提交
1560
      terrno = TSDB_CODE_OUT_OF_MEMORY;
L
Liu Jicong 已提交
1561 1562
      code = -1;
      goto END;
X
Xiaoyu Wang 已提交
1563
    }
1564

L
Liu Jicong 已提交
1565 1566 1567
    pWrapper->tmqRspType = TMQ_MSG_TYPE__EP_RSP;
    pWrapper->epoch = head->epoch;
    memcpy(&pWrapper->msg, pMsg->pData, sizeof(SMqRspHead));
L
Liu Jicong 已提交
1568
    tDecodeSMqAskEpRsp(POINTER_SHIFT(pMsg->pData, sizeof(SMqRspHead)), &pWrapper->msg);
L
Liu Jicong 已提交
1569

L
Liu Jicong 已提交
1570
    taosWriteQitem(tmq->mqueue, pWrapper);
1571
    tsem_post(&tmq->rspSem);
1572
  }
L
Liu Jicong 已提交
1573 1574

END:
1575 1576
  taosReleaseRef(tmqMgmt.rsetId, pParam->refId);

L
Liu Jicong 已提交
1577
  if (!async) {
L
Liu Jicong 已提交
1578
    tsem_post(&pParam->rspSem);
L
Liu Jicong 已提交
1579 1580
  } else {
    taosMemoryFree(pParam);
L
Liu Jicong 已提交
1581
  }
dengyihao's avatar
dengyihao 已提交
1582 1583

  taosMemoryFree(pMsg->pEpSet);
L
Liu Jicong 已提交
1584
  taosMemoryFree(pMsg->pData);
L
Liu Jicong 已提交
1585
  return code;
1586 1587
}

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

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

L
Liu Jicong 已提交
1605 1606
SMqMetaRspObj* tmqBuildMetaRspFromWrapper(SMqPollRspWrapper* pWrapper) {
  SMqMetaRspObj* pRspObj = taosMemoryCalloc(1, sizeof(SMqMetaRspObj));
L
Liu Jicong 已提交
1607
  pRspObj->resType = RES_TYPE__TMQ_META;
L
Liu Jicong 已提交
1608 1609 1610 1611 1612 1613 1614 1615
  tstrncpy(pRspObj->topic, pWrapper->topicHandle->topicName, TSDB_TOPIC_FNAME_LEN);
  tstrncpy(pRspObj->db, pWrapper->topicHandle->db, TSDB_DB_FNAME_LEN);
  pRspObj->vgId = pWrapper->vgHandle->vgId;

  memcpy(&pRspObj->metaRsp, &pWrapper->metaRsp, sizeof(SMqMetaRsp));
  return pRspObj;
}

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

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

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

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

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

1635 1636 1637 1638 1639
  // 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;
1640
    (*numOfRows) += rows;
1641 1642
  }

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

L
Liu Jicong 已提交
1646 1647
SMqTaosxRspObj* tmqBuildTaosxRspFromWrapper(SMqPollRspWrapper* pWrapper) {
  SMqTaosxRspObj* pRspObj = taosMemoryCalloc(1, sizeof(SMqTaosxRspObj));
1648
  pRspObj->resType = RES_TYPE__TMQ_METADATA;
L
Liu Jicong 已提交
1649 1650 1651 1652
  tstrncpy(pRspObj->topic, pWrapper->topicHandle->topicName, TSDB_TOPIC_FNAME_LEN);
  tstrncpy(pRspObj->db, pWrapper->topicHandle->db, TSDB_DB_FNAME_LEN);
  pRspObj->vgId = pWrapper->vgHandle->vgId;
  pRspObj->resIter = -1;
1653
  memcpy(&pRspObj->rsp, &pWrapper->taosxRsp, sizeof(STaosxRsp));
L
Liu Jicong 已提交
1654 1655 1656

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

  return pRspObj;
}

1664 1665 1666 1667 1668 1669 1670 1671 1672 1673 1674 1675 1676 1677 1678 1679 1680 1681 1682 1683 1684 1685 1686 1687 1688 1689 1690 1691 1692 1693 1694 1695 1696
static int32_t handleErrorBeforePoll(SMqClientVg* pVg, tmq_t* pTmq) {
  atomic_store_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE);
  tsem_post(&pTmq->rspSem);
  return -1;
}

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

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

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

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

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

  pParam->refId = pTmq->refId;
  pParam->epoch = pTmq->epoch;
X
Xiaoyu Wang 已提交
1697
  pParam->pVg = pVg;  // pVg may be released,fix it
1698 1699
  pParam->pTopic = pTopic;
  pParam->vgId = pVg->vgId;
H
Haojun Liao 已提交
1700
  pParam->requestId = req.reqId;
1701 1702 1703 1704 1705 1706 1707 1708 1709 1710 1711 1712 1713 1714 1715 1716 1717 1718 1719 1720 1721 1722 1723 1724

  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 已提交
1725
  tscDebug("consumer:0x%" PRIx64 " send poll to %s vgId:%d, epoch %d, req:%s, reqId:0x%" PRIx64,
1726 1727 1728 1729 1730 1731 1732 1733 1734
           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;
}

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

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

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

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

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

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

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

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

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

1806
    if (pRspWrapper == NULL) {
L
Liu Jicong 已提交
1807
      taosReadAllQitems(tmq->mqueue, tmq->qall);
1808
      taosGetQitem(tmq->qall, (void**)&pRspWrapper);
L
Liu Jicong 已提交
1809

1810
      if (pRspWrapper == NULL) {
L
Liu Jicong 已提交
1811 1812
        return NULL;
      }
X
Xiaoyu Wang 已提交
1813 1814
    }

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

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

1825
      int32_t consumerEpoch = atomic_load_32(&tmq->epoch);
1826 1827 1828 1829
      SMqDataRsp* pDataRsp = &pollRspWrapper->dataRsp;

      if (pDataRsp->head.epoch == consumerEpoch) {
        // todo fix it: race condition
L
Liu Jicong 已提交
1830
        SMqClientVg* pVg = pollRspWrapper->vgHandle;
1831 1832 1833 1834 1835 1836 1837 1838 1839 1840 1841

        // 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 已提交
1842
        atomic_store_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE);
H
Haojun Liao 已提交
1843

1844 1845 1846
        char buf[80];
        tFormatOffset(buf, 80, &pDataRsp->rspOffset);
        if (pDataRsp->blockNum == 0) {
H
Haojun Liao 已提交
1847 1848
          tscDebug("consumer:0x%" PRIx64 " empty block received, vgId:%d, offset:%s, vg total:%"PRId64" total:%"PRId64" reqId:0x%" PRIx64, tmq->consumerId,
                   pVg->vgId, buf, pVg->numOfRows, tmq->totalRows, pollRspWrapper->reqId);
1849
          pRspWrapper = tmqFreeRspWrapper(pRspWrapper);
L
Liu Jicong 已提交
1850
          taosFreeQitem(pollRspWrapper);
1851
        } else {  // build rsp
1852 1853
          int64_t numOfRows = 0;
          SMqRspObj* pRsp = tmqBuildRspFromWrapper(pollRspWrapper, pVg, &numOfRows);
H
Haojun Liao 已提交
1854 1855
          tmq->totalRows += numOfRows;

H
Haojun Liao 已提交
1856
          tscDebug("consumer:0x%" PRIx64 " process poll rsp, vgId:%d, offset:%s, blocks:%d, rows:%" PRId64
H
Haojun Liao 已提交
1857
                   " vg total:%" PRId64 " total:%" PRId64 ", reqId:0x%" PRIx64,
H
Haojun Liao 已提交
1858
                   tmq->consumerId, pVg->vgId, buf, pDataRsp->blockNum, numOfRows, pVg->numOfRows, tmq->totalRows,
H
Haojun Liao 已提交
1859
                   pollRspWrapper->reqId);
1860 1861 1862
          taosFreeQitem(pollRspWrapper);
          return pRsp;
        }
X
Xiaoyu Wang 已提交
1863
      } else {
H
Haojun Liao 已提交
1864 1865 1866
        SMqClientVg* pVg = pollRspWrapper->vgHandle;
        tscDebug("consumer:0x%" PRIx64 " vgId:%d msg discard since epoch mismatch: msg epoch %d, consumer epoch %d",
                 tmq->consumerId, pVg->vgId, pDataRsp->head.epoch, consumerEpoch);
1867
        pRspWrapper = tmqFreeRspWrapper(pRspWrapper);
L
Liu Jicong 已提交
1868 1869
        taosFreeQitem(pollRspWrapper);
      }
1870 1871
    } else if (pRspWrapper->tmqRspType == TMQ_MSG_TYPE__POLL_META_RSP) {
      SMqPollRspWrapper* pollRspWrapper = (SMqPollRspWrapper*)pRspWrapper;
L
Liu Jicong 已提交
1872
      int32_t            consumerEpoch = atomic_load_32(&tmq->epoch);
1873 1874 1875

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

L
Liu Jicong 已提交
1876
      if (pollRspWrapper->metaRsp.head.epoch == consumerEpoch) {
L
Liu Jicong 已提交
1877
        SMqClientVg* pVg = pollRspWrapper->vgHandle;
wmmhello's avatar
wmmhello 已提交
1878
        pVg->currentOffset = pollRspWrapper->metaRsp.rspOffset;
L
Liu Jicong 已提交
1879 1880
        atomic_store_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE);
        // build rsp
L
Liu Jicong 已提交
1881
        SMqMetaRspObj* pRsp = tmqBuildMetaRspFromWrapper(pollRspWrapper);
L
Liu Jicong 已提交
1882 1883 1884
        taosFreeQitem(pollRspWrapper);
        return pRsp;
      } else {
H
Haojun Liao 已提交
1885 1886
        tscDebug("consumer:0x%" PRIx64 " vgId:%d msg discard since epoch mismatch: msg epoch %d, consumer epoch %d",
                 tmq->consumerId, pollRspWrapper->vgHandle->vgId, pollRspWrapper->metaRsp.head.epoch, consumerEpoch);
1887
        pRspWrapper = tmqFreeRspWrapper(pRspWrapper);
L
Liu Jicong 已提交
1888
        taosFreeQitem(pollRspWrapper);
X
Xiaoyu Wang 已提交
1889
      }
1890 1891
    } else if (pRspWrapper->tmqRspType == TMQ_MSG_TYPE__TAOSX_RSP) {
      SMqPollRspWrapper* pollRspWrapper = (SMqPollRspWrapper*)pRspWrapper;
L
Liu Jicong 已提交
1892
      int32_t consumerEpoch = atomic_load_32(&tmq->epoch);
H
Haojun Liao 已提交
1893

L
Liu Jicong 已提交
1894 1895 1896 1897
      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 已提交
1898

L
Liu Jicong 已提交
1899
        if (pollRspWrapper->taosxRsp.blockNum == 0) {
H
Haojun Liao 已提交
1900 1901
          tscDebug("consumer:0x%" PRIx64 " taosx empty block received, vgId:%d, vg total:%" PRId64 " reqId:0x%" PRIx64,
                   tmq->consumerId, pVg->vgId, pVg->numOfRows, pollRspWrapper->reqId);
H
Haojun Liao 已提交
1902
          pVg->emptyBlockReceiveTs = taosGetTimestampMs();
1903
          pRspWrapper = tmqFreeRspWrapper(pRspWrapper);
H
Haojun Liao 已提交
1904
          taosFreeQitem(pollRspWrapper);
L
Liu Jicong 已提交
1905
          continue;
H
Haojun Liao 已提交
1906 1907
        } else {
          pVg->emptyBlockReceiveTs = 0; // reset the ts
L
Liu Jicong 已提交
1908
        }
wmmhello's avatar
wmmhello 已提交
1909

L
Liu Jicong 已提交
1910
        // build rsp
wmmhello's avatar
wmmhello 已提交
1911
        void* pRsp = NULL;
1912
        int64_t numOfRows = 0;
L
Liu Jicong 已提交
1913
        if (pollRspWrapper->taosxRsp.createTableNum == 0) {
1914
          pRsp = tmqBuildRspFromWrapper(pollRspWrapper, pVg, &numOfRows);
L
Liu Jicong 已提交
1915
        } else {
wmmhello's avatar
wmmhello 已提交
1916 1917
          pRsp = tmqBuildTaosxRspFromWrapper(pollRspWrapper);
        }
H
Haojun Liao 已提交
1918

1919 1920
        tmq->totalRows += numOfRows;

H
Haojun Liao 已提交
1921 1922
        char buf[80];
        tFormatOffset(buf, 80, &pVg->currentOffset);
H
Haojun Liao 已提交
1923
        tscDebug("consumer:0x%" PRIx64 " process taosx poll rsp, vgId:%d, offset:%s, blocks:%d, rows:%" PRId64
H
Haojun Liao 已提交
1924
                 ", vg total:%" PRId64 " total:%"PRId64" reqId:0x%" PRIx64,
H
Haojun Liao 已提交
1925
                 tmq->consumerId, pVg->vgId, buf, pollRspWrapper->dataRsp.blockNum, numOfRows, pVg->numOfRows,
H
Haojun Liao 已提交
1926
                 tmq->totalRows, pollRspWrapper->reqId);
H
Haojun Liao 已提交
1927 1928

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

L
Liu Jicong 已提交
1931
      } else {
H
Haojun Liao 已提交
1932 1933
        tscDebug("consumer:0x%" PRIx64 " vgId:%d msg discard since epoch mismatch: msg epoch %d, consumer epoch %d",
                 tmq->consumerId, pollRspWrapper->vgHandle->vgId, pollRspWrapper->taosxRsp.head.epoch, consumerEpoch);
1934
        pRspWrapper = tmqFreeRspWrapper(pRspWrapper);
L
Liu Jicong 已提交
1935 1936
        taosFreeQitem(pollRspWrapper);
      }
X
Xiaoyu Wang 已提交
1937
    } else {
H
Haojun Liao 已提交
1938 1939
      tscDebug("consumer:0x%" PRIx64 " not data msg received", tmq->consumerId);

X
Xiaoyu Wang 已提交
1940
      bool reset = false;
1941 1942
      tmqHandleNoPollRsp(tmq, pRspWrapper, &reset);
      taosFreeQitem(pRspWrapper);
X
Xiaoyu Wang 已提交
1943
      if (pollIfReset && reset) {
1944
        tscDebug("consumer:0x%" PRIx64 ", reset and repoll", tmq->consumerId);
1945
        tmqPollImpl(tmq, timeout);
X
Xiaoyu Wang 已提交
1946 1947 1948 1949 1950
      }
    }
  }
}

1951
TAOS_RES* tmq_consumer_poll(tmq_t* tmq, int64_t timeout) {
L
Liu Jicong 已提交
1952 1953
  void*   rspObj;
  int64_t startTime = taosGetTimestampMs();
L
Liu Jicong 已提交
1954

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

1957 1958 1959
#if 0
  tmqHandleAllDelayedTask(tmq);
  tmqPollImpl(tmq, timeout);
1960
  rspObj = tmqHandleAllRsp(tmq, timeout, false);
L
Liu Jicong 已提交
1961 1962
  if (rspObj) {
    return (TAOS_RES*)rspObj;
L
fix  
Liu Jicong 已提交
1963
  }
1964
#endif
X
Xiaoyu Wang 已提交
1965

1966
  // in no topic status, delayed task also need to be processed
L
Liu Jicong 已提交
1967
  if (atomic_load_8(&tmq->status) == TMQ_CONSUMER_STATUS__INIT) {
1968
    tscDebug("consumer:0x%" PRIx64 " poll return since consumer is init", tmq->consumerId);
1969
    taosMsleep(500);  //     sleep for a while
1970 1971 1972
    return NULL;
  }

wmmhello's avatar
wmmhello 已提交
1973
  while (atomic_load_8(&tmq->status) == TMQ_CONSUMER_STATUS__RECOVER) {
L
Liu Jicong 已提交
1974 1975
    int32_t retryCnt = 0;
    while (TSDB_CODE_MND_CONSUMER_NOT_READY == tmqAskEp(tmq, false)) {
H
Haojun Liao 已提交
1976
      if (retryCnt++ > 40) {
L
Liu Jicong 已提交
1977 1978
        return NULL;
      }
1979

H
Haojun Liao 已提交
1980
      tscDebug("consumer:0x%" PRIx64 " not ready, retry:%d/40 in 500ms", tmq->consumerId, retryCnt);
L
Liu Jicong 已提交
1981 1982 1983 1984
      taosMsleep(500);
    }
  }

X
Xiaoyu Wang 已提交
1985
  while (1) {
L
Liu Jicong 已提交
1986
    tmqHandleAllDelayedTask(tmq);
1987

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

1992
    rspObj = tmqHandleAllRsp(tmq, timeout, false);
L
Liu Jicong 已提交
1993
    if (rspObj) {
1994
      tscDebug("consumer:0x%" PRIx64 " return rsp %p", tmq->consumerId, rspObj);
L
Liu Jicong 已提交
1995
      return (TAOS_RES*)rspObj;
L
Liu Jicong 已提交
1996
    } else if (terrno == TSDB_CODE_TQ_NO_COMMITTED_OFFSET) {
1997
      tscDebug("consumer:0x%" PRIx64 " return null since no committed offset", tmq->consumerId);
L
Liu Jicong 已提交
1998
      return NULL;
X
Xiaoyu Wang 已提交
1999
    }
2000

2001
    if (timeout >= 0) {
L
Liu Jicong 已提交
2002
      int64_t currentTime = taosGetTimestampMs();
2003 2004 2005
      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 已提交
2006
                 tmq->consumerId, tmq->epoch, startTime, currentTime);
X
Xiaoyu Wang 已提交
2007 2008
        return NULL;
      }
2009
      tsem_timewait(&tmq->rspSem, (timeout - elapsedTime));
L
Liu Jicong 已提交
2010 2011
    } else {
      // use tsem_timewait instead of tsem_wait to avoid unexpected stuck
L
Liu Jicong 已提交
2012
      tsem_timewait(&tmq->rspSem, 1000);
X
Xiaoyu Wang 已提交
2013 2014 2015 2016
    }
  }
}

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

2020
  if (tmq->status == TMQ_CONSUMER_STATUS__READY) {
L
Liu Jicong 已提交
2021 2022
    int32_t rsp = tmq_commit_sync(tmq, NULL);
    if (rsp != 0) {
L
Liu Jicong 已提交
2023
      return rsp;
2024 2025
    }

2026 2027 2028 2029 2030 2031 2032 2033 2034 2035 2036 2037 2038 2039 2040 2041 2042 2043
    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 已提交
2044
    int32_t     retryCnt = 0;
2045
    tmq_list_t* lst = tmq_list_new();
L
Liu Jicong 已提交
2046 2047 2048 2049 2050 2051 2052 2053 2054 2055
    while (1) {
      rsp = tmq_subscribe(tmq, lst);
      if (rsp != TSDB_CODE_MND_CONSUMER_NOT_READY || retryCnt > 5) {
        break;
      } else {
        retryCnt++;
        taosMsleep(500);
      }
    }

2056
    tmq_list_destroy(lst);
L
Liu Jicong 已提交
2057
  }
H
Haojun Liao 已提交
2058

2059
  taosRemoveRef(tmqMgmt.rsetId, tmq->refId);
L
Liu Jicong 已提交
2060
  return 0;
2061
}
L
Liu Jicong 已提交
2062

L
Liu Jicong 已提交
2063 2064
const char* tmq_err2str(int32_t err) {
  if (err == 0) {
L
Liu Jicong 已提交
2065
    return "success";
L
Liu Jicong 已提交
2066
  } else if (err == -1) {
L
Liu Jicong 已提交
2067 2068 2069
    return "fail";
  } else {
    return tstrerror(err);
L
Liu Jicong 已提交
2070 2071
  }
}
L
Liu Jicong 已提交
2072

L
Liu Jicong 已提交
2073 2074 2075 2076 2077
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;
2078 2079
  } else if (TD_RES_TMQ_METADATA(res)) {
    return TMQ_RES_METADATA;
L
Liu Jicong 已提交
2080 2081 2082 2083 2084
  } else {
    return TMQ_RES_INVALID;
  }
}

L
Liu Jicong 已提交
2085
const char* tmq_get_topic_name(TAOS_RES* res) {
L
Liu Jicong 已提交
2086 2087
  if (TD_RES_TMQ(res)) {
    SMqRspObj* pRspObj = (SMqRspObj*)res;
L
Liu Jicong 已提交
2088
    return strchr(pRspObj->topic, '.') + 1;
L
Liu Jicong 已提交
2089 2090 2091
  } else if (TD_RES_TMQ_META(res)) {
    SMqMetaRspObj* pMetaRspObj = (SMqMetaRspObj*)res;
    return strchr(pMetaRspObj->topic, '.') + 1;
2092 2093 2094
  } else if (TD_RES_TMQ_METADATA(res)) {
    SMqTaosxRspObj* pRspObj = (SMqTaosxRspObj*)res;
    return strchr(pRspObj->topic, '.') + 1;
L
Liu Jicong 已提交
2095 2096 2097 2098 2099
  } else {
    return NULL;
  }
}

L
Liu Jicong 已提交
2100 2101 2102 2103
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 已提交
2104 2105 2106
  } else if (TD_RES_TMQ_META(res)) {
    SMqMetaRspObj* pMetaRspObj = (SMqMetaRspObj*)res;
    return strchr(pMetaRspObj->db, '.') + 1;
2107 2108 2109
  } else if (TD_RES_TMQ_METADATA(res)) {
    SMqTaosxRspObj* pRspObj = (SMqTaosxRspObj*)res;
    return strchr(pRspObj->db, '.') + 1;
L
Liu Jicong 已提交
2110 2111 2112 2113 2114
  } else {
    return NULL;
  }
}

L
Liu Jicong 已提交
2115 2116 2117 2118
int32_t tmq_get_vgroup_id(TAOS_RES* res) {
  if (TD_RES_TMQ(res)) {
    SMqRspObj* pRspObj = (SMqRspObj*)res;
    return pRspObj->vgId;
L
Liu Jicong 已提交
2119 2120 2121
  } else if (TD_RES_TMQ_META(res)) {
    SMqMetaRspObj* pMetaRspObj = (SMqMetaRspObj*)res;
    return pMetaRspObj->vgId;
2122
  } else if (TD_RES_TMQ_METADATA(res)) {
2123 2124
    SMqTaosxRspObj* pRspObj = (SMqTaosxRspObj*)res;
    return pRspObj->vgId;
L
Liu Jicong 已提交
2125 2126 2127 2128
  } else {
    return -1;
  }
}
L
Liu Jicong 已提交
2129 2130 2131 2132 2133 2134 2135 2136

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;
    }
2137
    return (const char*)taosArrayGetP(pRspObj->rsp.blockTbName, pRspObj->resIter);
2138 2139
  } else if (TD_RES_TMQ_METADATA(res)) {
    SMqTaosxRspObj* pRspObj = (SMqTaosxRspObj*)res;
L
Liu Jicong 已提交
2140 2141 2142
    if (!pRspObj->rsp.withTbName || pRspObj->rsp.blockTbName == NULL || pRspObj->resIter < 0 ||
        pRspObj->resIter >= pRspObj->rsp.blockNum) {
      return NULL;
2143
    }
L
Liu Jicong 已提交
2144 2145
    return (const char*)taosArrayGetP(pRspObj->rsp.blockTbName, pRspObj->resIter);
  }
L
Liu Jicong 已提交
2146 2147
  return NULL;
}
2148

L
Liu Jicong 已提交
2149
void tmq_commit_async(tmq_t* tmq, const TAOS_RES* msg, tmq_commit_cb* cb, void* param) {
L
Liu Jicong 已提交
2150
  tmqCommitInner(tmq, msg, 0, 1, cb, param);
L
Liu Jicong 已提交
2151 2152
}

2153
int32_t tmq_commit_sync(tmq_t* tmq, const TAOS_RES* msg) {
L
Liu Jicong 已提交
2154
  return tmqCommitInner(tmq, msg, 0, 0, NULL, NULL);
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 2236 2237 2238 2239 2240 2241 2242 2243 2244 2245

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) {
2246 2247 2248
  int64_t refId = pParamSet->refId;

  tmq_t* tmq = taosAcquireRef(tmqMgmt.rsetId, refId);
2249 2250 2251 2252 2253 2254 2255 2256 2257 2258 2259 2260 2261 2262
  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);
2263
    } else if (!pParamSet->automatic && pParamSet->userCb) { // sem post
2264 2265
      pParamSet->userCb(tmq, pParamSet->rspErr, pParamSet->userParam);
    }
2266

2267 2268 2269 2270 2271 2272 2273 2274 2275
    taosMemoryFree(pParamSet);
  } else {
    tsem_post(&pParamSet->rspSem);
  }

#if 0
  taosArrayDestroyP(pParamSet->successfulOffsets, taosMemoryFree);
    taosArrayDestroyP(pParamSet->failedOffsets, taosMemoryFree);
#endif
2276 2277

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

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 已提交
2284 2285
    tscDebug("consumer:0x%" PRIx64 " topic:%s vgId:%d all commit-rsp received, commit completed", consumerId, pTopic,
             vgId);
2286
    tmqCommitDone(pParamSet);
H
Haojun Liao 已提交
2287 2288 2289
  } else {
    tscDebug("consumer:0x%" PRIx64 " topic:%s vgId:%d commit-rsp received, remain:%d", consumerId, pTopic, vgId,
             waitingRspNum);
2290 2291
  }
}