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

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

X
Xiaoyu Wang 已提交
27
struct SMqMgmt {
28 29 30
  int8_t  inited;
  tmr_h   timer;
  int32_t rsetId;
31
};
L
Liu Jicong 已提交
32

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

L
Liu Jicong 已提交
37 38 39 40 41 42
typedef struct {
  int8_t  tmqRspType;
  int32_t epoch;
} SMqRspWrapper;

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

L
Liu Jicong 已提交
48
struct tmq_list_t {
L
Liu Jicong 已提交
49
  SArray container;
L
Liu Jicong 已提交
50
};
L
Liu Jicong 已提交
51

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

struct tmq_t {
71
  int64_t refId;
L
Liu Jicong 已提交
72
  // conf
73 74 75 76 77 78 79 80 81
  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;
82

L
Liu Jicong 已提交
83 84
  tmq_commit_cb* commitCb;
  void*          commitCbUserParam;
L
Liu Jicong 已提交
85 86 87 88

  // status
  int8_t  status;
  int32_t epoch;
L
Liu Jicong 已提交
89 90
#if 0
  int8_t  epStatus;
L
Liu Jicong 已提交
91
  int32_t epSkipCnt;
L
Liu Jicong 已提交
92
#endif
L
Liu Jicong 已提交
93 94
  int64_t pollCnt;

L
Liu Jicong 已提交
95
  // timer
96 97
  tmr_h hbLiveTimer;
  tmr_h epTimer;
L
Liu Jicong 已提交
98 99 100
  tmr_h reportTimer;
  tmr_h commitTimer;

H
Haojun Liao 已提交
101 102 103 104 105 106 107
  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 已提交
108 109
};

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

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

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

L
Liu Jicong 已提交
128
typedef struct {
129 130 131
  // statistics
  int64_t pollCnt;
  // offset
L
Liu Jicong 已提交
132 133
  STqOffsetVal committedOffset;
  STqOffsetVal currentOffset;
L
Liu Jicong 已提交
134
  // connection info
135
  int32_t vgId;
X
Xiaoyu Wang 已提交
136
  int32_t vgStatus;
L
Liu Jicong 已提交
137
  int32_t vgSkipCnt;
138 139 140
  SEpSet  epSet;
} SMqClientVg;

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

L
Liu Jicong 已提交
148 149 150 151 152
typedef struct {
  int8_t          tmqRspType;
  int32_t         epoch;
  SMqClientVg*    vgHandle;
  SMqClientTopic* topicHandle;
L
Liu Jicong 已提交
153
  union {
L
Liu Jicong 已提交
154 155
    SMqDataRsp dataRsp;
    SMqMetaRsp metaRsp;
L
Liu Jicong 已提交
156
    STaosxRsp  taosxRsp;
L
Liu Jicong 已提交
157
  };
L
Liu Jicong 已提交
158 159
} SMqPollRspWrapper;

L
Liu Jicong 已提交
160
typedef struct {
161 162
  int64_t refId;
  int32_t epoch;
L
Liu Jicong 已提交
163 164
  tsem_t  rspSem;
  int32_t rspErr;
L
Liu Jicong 已提交
165
} SMqSubscribeCbParam;
L
Liu Jicong 已提交
166

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

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

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

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

208
static int32_t tmqAskEp(tmq_t* tmq, bool async);
209 210
static int32_t makeTopicVgroupKey(char* dst, const char* topicName, int32_t vg);
static int32_t tmqCommitDone(SMqCommitCbParamSet* pParamSet);
211
static void tmqCommitRspCountDown(SMqCommitCbParamSet* pParamSet, int64_t consumerId, const char* pTopic, int32_t vgId);
212

213
tmq_conf_t* tmq_conf_new() {
wafwerar's avatar
wafwerar 已提交
214
  tmq_conf_t* conf = taosMemoryCalloc(1, sizeof(tmq_conf_t));
215 216 217 218 219
  if (conf == NULL) {
    terrno = TSDB_CODE_OUT_OF_MEMORY;
    return conf;
  }

220
  conf->withTbName = false;
L
Liu Jicong 已提交
221
  conf->autoCommit = true;
L
Liu Jicong 已提交
222
  conf->autoCommitInterval = 5000;
223
  conf->resetOffset = TMQ_OFFSET__RESET_EARLIEAST;
224
  conf->hbBgEnable = true;
225

226 227 228
  return conf;
}

L
Liu Jicong 已提交
229
void tmq_conf_destroy(tmq_conf_t* conf) {
L
Liu Jicong 已提交
230
  if (conf) {
231 232 233 234 235 236 237 238 239
    if (conf->ip) {
      taosMemoryFree(conf->ip);
    }
    if (conf->user) {
      taosMemoryFree(conf->user);
    }
    if (conf->pass) {
      taosMemoryFree(conf->pass);
    }
L
Liu Jicong 已提交
240 241
    taosMemoryFree(conf);
  }
L
Liu Jicong 已提交
242 243 244
}

tmq_conf_res_t tmq_conf_set(tmq_conf_t* conf, const char* key, const char* value) {
245
  if (strcasecmp(key, "group.id") == 0) {
L
Liu Jicong 已提交
246
    tstrncpy(conf->groupId, value, TSDB_CGROUP_LEN);
L
Liu Jicong 已提交
247
    return TMQ_CONF_OK;
248
  }
L
Liu Jicong 已提交
249

250
  if (strcasecmp(key, "client.id") == 0) {
L
Liu Jicong 已提交
251
    tstrncpy(conf->clientId, value, 256);
L
Liu Jicong 已提交
252 253
    return TMQ_CONF_OK;
  }
L
Liu Jicong 已提交
254

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

267
  if (strcasecmp(key, "auto.commit.interval.ms") == 0) {
L
Liu Jicong 已提交
268 269 270 271
    conf->autoCommitInterval = atoi(value);
    return TMQ_CONF_OK;
  }

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

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

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

311
  if (strcasecmp(key, "experimental.snapshot.batch.size") == 0) {
L
Liu Jicong 已提交
312 313 314 315
    conf->snapBatchSize = atoi(value);
    return TMQ_CONF_OK;
  }

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

328
  if (strcasecmp(key, "td.connect.ip") == 0) {
329
    conf->ip = taosStrdup(value);
L
Liu Jicong 已提交
330 331
    return TMQ_CONF_OK;
  }
332
  if (strcasecmp(key, "td.connect.user") == 0) {
333
    conf->user = taosStrdup(value);
L
Liu Jicong 已提交
334 335
    return TMQ_CONF_OK;
  }
336
  if (strcasecmp(key, "td.connect.pass") == 0) {
337
    conf->pass = taosStrdup(value);
L
Liu Jicong 已提交
338 339
    return TMQ_CONF_OK;
  }
340
  if (strcasecmp(key, "td.connect.port") == 0) {
L
Liu Jicong 已提交
341 342 343
    conf->port = atoi(value);
    return TMQ_CONF_OK;
  }
344
  if (strcasecmp(key, "td.connect.db") == 0) {
L
Liu Jicong 已提交
345 346 347
    return TMQ_CONF_OK;
  }

L
Liu Jicong 已提交
348
  return TMQ_CONF_UNKNOWN;
349 350 351
}

tmq_list_t* tmq_list_new() {
L
Liu Jicong 已提交
352
  return (tmq_list_t*)taosArrayInit(0, sizeof(void*));
353 354
}

L
Liu Jicong 已提交
355 356
int32_t tmq_list_append(tmq_list_t* list, const char* src) {
  SArray* container = &list->container;
357
  if (src == NULL || src[0] == 0) return -1;
358
  char* topic = taosStrdup(src);
359 360 361
  if (topic[0] != '`') {
    strtolower(topic, src);
  }
L
fix  
Liu Jicong 已提交
362
  if (taosArrayPush(container, &topic) == NULL) return -1;
363 364 365
  return 0;
}

L
Liu Jicong 已提交
366
void tmq_list_destroy(tmq_list_t* list) {
L
Liu Jicong 已提交
367
  SArray* container = &list->container;
L
Liu Jicong 已提交
368
  taosArrayDestroyP(container, taosMemoryFree);
L
Liu Jicong 已提交
369 370
}

L
Liu Jicong 已提交
371 372 373 374 375 376 377 378 379 380
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;
}

381 382 383 384 385 386 387 388 389 390 391 392 393 394 395 396 397 398 399 400 401 402
static void updateVgEpset(tmq_t* pTmq, SMqCommitCbParam* pParam, SEpSet* pEpSet) {
  int32_t numOfTopics = taosArrayGetSize(pTmq->clientTopics);
  for(int32_t i = 0; i < numOfTopics; ++i) {
    SMqClientTopic* pTopic = taosArrayGet(pTmq->clientTopics, i);
    if (strcmp(pTopic->topicName, pParam->topicName) != 0) {
      continue;
    }

    int32_t numOfVgs = taosArrayGetSize(pTopic->vgs);
    for(int32_t j = 0; j < numOfVgs; ++j) {
      SMqClientVg* pClientVg = taosArrayGet(pTopic->vgs, j);
      if (pClientVg->vgId == pParam->vgId) {
        SEp* pEp = GET_ACTIVE_EP(pEpSet);
        SEp* pOld = GET_ACTIVE_EP(&(pClientVg->epSet));
        uDebug("subKey:%s update the epset vgId:%d, ep:%s:%d, old ep:%s:%d", pParam->pOffset->subKey, pParam->vgId,
               pEp->fqdn, pEp->port, pOld->fqdn, pOld->port);
        pClientVg->epSet = *pEpSet;
        break;
      }
    }

    break;
L
Liu Jicong 已提交
403 404
  }
}
405
// todo retry to send the commit if failed
H
Haojun Liao 已提交
406
static int32_t tmqCommitCb(void* param, SDataBuf* pBuf, int32_t code) {
407
  SMqCommitCbParam*    pParam = (SMqCommitCbParam*)param;
408
  SMqCommitCbParamSet* pParamSet = (SMqCommitCbParamSet*)pParam->params;
H
Haojun Liao 已提交
409

410
  // push into array
L
Liu Jicong 已提交
411
#if 0
412 413 414 415 416
  if (code == 0) {
    taosArrayPush(pParamSet->failedOffsets, &pParam->pOffset);
  } else {
    taosArrayPush(pParamSet->successfulOffsets, &pParam->pOffset);
  }
L
Liu Jicong 已提交
417
#endif
L
Liu Jicong 已提交
418

419
  // update the epset if needed
H
Haojun Liao 已提交
420 421
  if (pBuf->pEpSet != NULL) {
    taosThreadMutexLock(&pParam->pTmq->lock);
422
    updateVgEpset(pParam->pTmq, pParam, pBuf->pEpSet);
H
Haojun Liao 已提交
423
    taosThreadMutexUnlock(&pParam->pTmq->lock);
424 425
  }

L
Liu Jicong 已提交
426
  taosMemoryFree(pParam->pOffset);
L
Liu Jicong 已提交
427
  taosMemoryFree(pBuf->pData);
dengyihao's avatar
dengyihao 已提交
428
  taosMemoryFree(pBuf->pEpSet);
L
Liu Jicong 已提交
429

430
  tmqCommitRspCountDown(pParamSet, pParam->pTmq->consumerId, pParam->topicName, pParam->vgId);
431 432 433
  return 0;
}

H
Haojun Liao 已提交
434
static int32_t tmqSendCommitReq(tmq_t* tmq, SMqClientVg* pVg, const char* pTopicName, SMqCommitCbParamSet* pParamSet) {
L
Liu Jicong 已提交
435 436 437 438 439
  STqOffset* pOffset = taosMemoryCalloc(1, sizeof(STqOffset));
  if (pOffset == NULL) {
    terrno = TSDB_CODE_OUT_OF_MEMORY;
    return -1;
  }
440

L
Liu Jicong 已提交
441
  pOffset->val = pVg->currentOffset;
442

L
Liu Jicong 已提交
443 444 445
  int32_t groupLen = strlen(tmq->groupId);
  memcpy(pOffset->subKey, tmq->groupId, groupLen);
  pOffset->subKey[groupLen] = TMQ_SEPARATOR;
H
Haojun Liao 已提交
446
  strcpy(pOffset->subKey + groupLen + 1, pTopicName);
L
Liu Jicong 已提交
447

L
Liu Jicong 已提交
448 449 450 451 452 453
  int32_t len;
  int32_t code;
  tEncodeSize(tEncodeSTqOffset, pOffset, len, code);
  if (code < 0) {
    return -1;
  }
454

L
Liu Jicong 已提交
455
  void* buf = taosMemoryCalloc(1, sizeof(SMsgHead) + len);
L
Liu Jicong 已提交
456 457 458 459
  if (buf == NULL) {
    taosMemoryFree(pOffset);
    return -1;
  }
460

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

L
Liu Jicong 已提交
463 464 465 466 467
  void* abuf = POINTER_SHIFT(buf, sizeof(SMsgHead));

  SEncoder encoder;
  tEncoderInit(&encoder, abuf, len);
  tEncodeSTqOffset(&encoder, pOffset);
L
Liu Jicong 已提交
468
  tEncoderClear(&encoder);
L
Liu Jicong 已提交
469 470

  // build param
471
  SMqCommitCbParam* pParam = taosMemoryCalloc(1, sizeof(SMqCommitCbParam));
L
Liu Jicong 已提交
472
  if (pParam == NULL) {
L
Liu Jicong 已提交
473
    taosMemoryFree(pOffset);
L
Liu Jicong 已提交
474 475 476
    taosMemoryFree(buf);
    return -1;
  }
477

L
Liu Jicong 已提交
478 479
  pParam->params = pParamSet;
  pParam->pOffset = pOffset;
H
Haojun Liao 已提交
480 481 482
  pParam->vgId = pVg->vgId;
  pParam->pTmq = tmq;

H
Haojun Liao 已提交
483
  tstrncpy(pParam->topicName, pTopicName, tListLen(pParam->topicName));
L
Liu Jicong 已提交
484 485 486 487

  // build send info
  SMsgSendInfo* pMsgSendInfo = taosMemoryCalloc(1, sizeof(SMsgSendInfo));
  if (pMsgSendInfo == NULL) {
L
Liu Jicong 已提交
488
    taosMemoryFree(pOffset);
L
Liu Jicong 已提交
489 490
    taosMemoryFree(buf);
    taosMemoryFree(pParam);
L
Liu Jicong 已提交
491 492
    return -1;
  }
493

L
Liu Jicong 已提交
494 495 496 497 498 499
  pMsgSendInfo->msgInfo = (SDataBuf){
      .pData = buf,
      .len = sizeof(SMsgHead) + len,
      .handle = NULL,
  };

H
Haojun Liao 已提交
500
  SEp* pEp = GET_ACTIVE_EP(&pVg->epSet);
X
Xiaoyu Wang 已提交
501 502
  tscDebug("consumer:0x%" PRIx64 " topic:%s on vgId:%d offset:%" PRId64 " prev:%" PRId64 ", ep:%s:%d", tmq->consumerId,
           pOffset->subKey, pVg->vgId, pOffset->val.version, pVg->committedOffset.version, pEp->fqdn, pEp->port);
L
Liu Jicong 已提交
503

504
  // TODO: put into cb, the commit offset should be move to the callback function
L
Liu Jicong 已提交
505
  pVg->committedOffset = pVg->currentOffset;
L
Liu Jicong 已提交
506 507 508 509

  pMsgSendInfo->requestId = generateRequestId();
  pMsgSendInfo->requestObjRefId = 0;
  pMsgSendInfo->param = pParam;
L
Liu Jicong 已提交
510
  pMsgSendInfo->paramFreeFp = taosMemoryFree;
511
  pMsgSendInfo->fp = tmqCommitCb;
L
Liu Jicong 已提交
512
  pMsgSendInfo->msgType = TDMT_VND_TMQ_COMMIT_OFFSET;
L
Liu Jicong 已提交
513

L
Liu Jicong 已提交
514 515 516
  atomic_add_fetch_32(&pParamSet->waitingRspNum, 1);
  atomic_add_fetch_32(&pParamSet->totalRspNum, 1);

L
Liu Jicong 已提交
517 518 519 520 521
  int64_t transporterId = 0;
  asyncSendMsgToServer(tmq->pTscObj->pAppInfo->pTransporter, &pVg->epSet, &transporterId, pMsgSendInfo);
  return 0;
}

H
Haojun Liao 已提交
522
static int32_t tmqCommitMsgImpl(tmq_t* tmq, const TAOS_RES* msg, int8_t async, tmq_commit_cb* userCb, void* userParam) {
L
Liu Jicong 已提交
523 524 525 526 527 528 529 530 531 532
  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 已提交
533
  } else if (TD_RES_TMQ_METADATA(msg)) {
534 535 536
    SMqTaosxRspObj* pRspObj = (SMqTaosxRspObj*)msg;
    topic = pRspObj->topic;
    vgId = pRspObj->vgId;
L
Liu Jicong 已提交
537 538 539 540 541 542 543 544 545
  } 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 已提交
546

547 548
  pParamSet->refId = tmq->refId;
  pParamSet->epoch = tmq->epoch;
L
Liu Jicong 已提交
549 550 551 552 553 554
  pParamSet->automatic = 0;
  pParamSet->async = async;
  pParamSet->userCb = userCb;
  pParamSet->userParam = userParam;
  tsem_init(&pParamSet->rspSem, 0, 0);

L
Liu Jicong 已提交
555 556
  int32_t code = -1;

H
Haojun Liao 已提交
557
  taosThreadMutexLock(&tmq->lock);
L
Liu Jicong 已提交
558 559 560 561 562
  for (int32_t i = 0; i < taosArrayGetSize(tmq->clientTopics); i++) {
    SMqClientTopic* pTopic = taosArrayGet(tmq->clientTopics, i);
    if (strcmp(pTopic->topicName, topic) != 0) continue;
    for (int32_t j = 0; j < taosArrayGetSize(pTopic->vgs); j++) {
      SMqClientVg* pVg = taosArrayGet(pTopic->vgs, j);
H
Haojun Liao 已提交
563 564 565
      if (pVg->vgId != vgId) {
        continue;
      }
L
Liu Jicong 已提交
566

L
Liu Jicong 已提交
567
      if (pVg->currentOffset.type > 0 && !tOffsetEqual(&pVg->currentOffset, &pVg->committedOffset)) {
H
Haojun Liao 已提交
568
        if (tmqSendCommitReq(tmq, pVg, pTopic->topicName, pParamSet) < 0) {
L
Liu Jicong 已提交
569 570
          tsem_destroy(&pParamSet->rspSem);
          taosMemoryFree(pParamSet);
L
Liu Jicong 已提交
571
          goto FAIL;
L
Liu Jicong 已提交
572
        }
L
Liu Jicong 已提交
573
        goto HANDLE_RSP;
L
Liu Jicong 已提交
574 575
      }
    }
L
Liu Jicong 已提交
576
  }
L
Liu Jicong 已提交
577

L
Liu Jicong 已提交
578 579 580 581
HANDLE_RSP:
  if (pParamSet->totalRspNum == 0) {
    tsem_destroy(&pParamSet->rspSem);
    taosMemoryFree(pParamSet);
H
Haojun Liao 已提交
582
    taosThreadMutexUnlock(&tmq->lock);
L
Liu Jicong 已提交
583 584 585
    return 0;
  }

L
Liu Jicong 已提交
586 587 588 589
  if (!async) {
    tsem_wait(&pParamSet->rspSem);
    code = pParamSet->rspErr;
    tsem_destroy(&pParamSet->rspSem);
L
Liu Jicong 已提交
590
    taosMemoryFree(pParamSet);
H
Haojun Liao 已提交
591
    taosThreadMutexUnlock(&tmq->lock);
L
Liu Jicong 已提交
592 593 594 595 596 597
    return code;
  } else {
    code = 0;
  }

FAIL:
H
Haojun Liao 已提交
598
  taosThreadMutexUnlock(&tmq->lock);
L
Liu Jicong 已提交
599 600 601
  if (code != 0 && async) {
    userCb(tmq, code, userParam);
  }
H
Haojun Liao 已提交
602

L
Liu Jicong 已提交
603 604 605
  return 0;
}

L
Liu Jicong 已提交
606 607
static int32_t tmqCommitConsumerImpl(tmq_t* tmq, int8_t automatic, int8_t async, tmq_commit_cb* userCb,
                                     void* userParam) {
L
Liu Jicong 已提交
608 609
  int32_t code = -1;

610 611
  SMqCommitCbParamSet* pParamSet = taosMemoryCalloc(1, sizeof(SMqCommitCbParamSet));
  if (pParamSet == NULL) {
L
Liu Jicong 已提交
612 613 614 615 616 617 618 619
    code = TSDB_CODE_OUT_OF_MEMORY;
    if (async) {
      if (automatic) {
        tmq->commitCb(tmq, code, tmq->commitCbUserParam);
      } else {
        userCb(tmq, code, userParam);
      }
    }
620 621
    return -1;
  }
622 623 624 625

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

626 627 628 629 630 631
  pParamSet->automatic = automatic;
  pParamSet->async = async;
  pParamSet->userCb = userCb;
  pParamSet->userParam = userParam;
  tsem_init(&pParamSet->rspSem, 0, 0);

632 633 634
  // init as 1 to prevent concurrency issue
  pParamSet->waitingRspNum = 1;

H
Haojun Liao 已提交
635
  taosThreadMutexLock(&tmq->lock);
636
  int32_t numOfTopics = taosArrayGetSize(tmq->clientTopics);
H
Haojun Liao 已提交
637

X
Xiaoyu Wang 已提交
638
  tscDebug("consumer:0x%" PRIx64 " start to commit offset for %d topics", tmq->consumerId, numOfTopics);
639 640

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

644 645
    tscDebug("consumer:0x%" PRIx64 " commit offset for topics:%s, numOfVgs:%d", tmq->consumerId, pTopic->topicName,
             numOfVgroups);
646
    for (int32_t j = 0; j < numOfVgroups; j++) {
H
Haojun Liao 已提交
647 648 649
      SMqClientVg clientVg = *(SMqClientVg*)taosArrayGet(pTopic->vgs, j);
      if (clientVg.currentOffset.type > 0 && !tOffsetEqual(&clientVg.currentOffset, &clientVg.committedOffset)) {
        if (tmqSendCommitReq(tmq, &clientVg, pTopic->topicName, pParamSet) < 0) {
L
Liu Jicong 已提交
650 651
          continue;
        }
652 653
      } else {
        tscDebug("consumer:0x%" PRIx64 " topic:%s vgId:%d, not commit, current:%" PRId64 ", ordinal:%d/%d",
H
Haojun Liao 已提交
654
                 tmq->consumerId, pTopic->topicName, clientVg.vgId, clientVg.currentOffset.version, j + 1, numOfVgroups);
655 656 657 658
      }
    }
  }

659 660
  tscDebug("consumer:0x%" PRIx64 " total commit:%d for %d topics", tmq->consumerId, pParamSet->waitingRspNum,
           numOfTopics);
H
Haojun Liao 已提交
661 662
  taosThreadMutexUnlock(&tmq->lock);

L
Liu Jicong 已提交
663
  // no request is sent
L
Liu Jicong 已提交
664 665 666 667 668 669
  if (pParamSet->totalRspNum == 0) {
    tsem_destroy(&pParamSet->rspSem);
    taosMemoryFree(pParamSet);
    return 0;
  }

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

673 674 675 676
  if (!async) {
    tsem_wait(&pParamSet->rspSem);
    code = pParamSet->rspErr;
    tsem_destroy(&pParamSet->rspSem);
677
    taosMemoryFree(pParamSet);
L
Liu Jicong 已提交
678
#if 0
679 680
    taosArrayDestroyP(pParamSet->successfulOffsets, taosMemoryFree);
    taosArrayDestroyP(pParamSet->failedOffsets, taosMemoryFree);
L
Liu Jicong 已提交
681
#endif
L
Liu Jicong 已提交
682
  }
683

L
Liu Jicong 已提交
684 685 686
  return code;
}

687 688
static int32_t tmqCommitInner(tmq_t* tmq, const TAOS_RES* msg, int8_t automatic, int8_t async, tmq_commit_cb* userCb,
                              void* userParam) {
689
  if (msg) { // user invoked commit?
L
Liu Jicong 已提交
690
    return tmqCommitMsgImpl(tmq, msg, async, userCb, userParam);
691
  } else {  // this for auto commit
L
Liu Jicong 已提交
692 693
    return tmqCommitConsumerImpl(tmq, automatic, async, userCb, userParam);
  }
694 695
}

696
void tmqAssignAskEpTask(void* param, void* tmrId) {
697 698 699
  int64_t refId = *(int64_t*)param;
  tmq_t*  tmq = taosAcquireRef(tmqMgmt.rsetId, refId);
  if (tmq != NULL) {
S
Shengliang Guan 已提交
700
    int8_t* pTaskType = taosAllocateQitem(sizeof(int8_t), DEF_QITEM, 0);
701 702 703 704 705
    *pTaskType = TMQ_DELAYED_TASK__ASK_EP;
    taosWriteQitem(tmq->delayedTask, pTaskType);
    tsem_post(&tmq->rspSem);
  }
  taosMemoryFree(param);
L
Liu Jicong 已提交
706 707 708
}

void tmqAssignDelayedCommitTask(void* param, void* tmrId) {
709 710 711
  int64_t refId = *(int64_t*)param;
  tmq_t*  tmq = taosAcquireRef(tmqMgmt.rsetId, refId);
  if (tmq != NULL) {
S
Shengliang Guan 已提交
712
    int8_t* pTaskType = taosAllocateQitem(sizeof(int8_t), DEF_QITEM, 0);
713 714 715 716 717
    *pTaskType = TMQ_DELAYED_TASK__COMMIT;
    taosWriteQitem(tmq->delayedTask, pTaskType);
    tsem_post(&tmq->rspSem);
  }
  taosMemoryFree(param);
L
Liu Jicong 已提交
718 719 720
}

void tmqAssignDelayedReportTask(void* param, void* tmrId) {
721 722 723
  int64_t refId = *(int64_t*)param;
  tmq_t*  tmq = taosAcquireRef(tmqMgmt.rsetId, refId);
  if (tmq != NULL) {
S
Shengliang Guan 已提交
724
    int8_t* pTaskType = taosAllocateQitem(sizeof(int8_t), DEF_QITEM, 0);
725 726 727 728 729
    *pTaskType = TMQ_DELAYED_TASK__REPORT;
    taosWriteQitem(tmq->delayedTask, pTaskType);
    tsem_post(&tmq->rspSem);
  }
  taosMemoryFree(param);
L
Liu Jicong 已提交
730 731
}

732
int32_t tmqHbCb(void* param, SDataBuf* pMsg, int32_t code) {
dengyihao's avatar
dengyihao 已提交
733 734 735 736
  if (pMsg) {
    taosMemoryFree(pMsg->pData);
    taosMemoryFree(pMsg->pEpSet);
  }
737 738 739 740
  return 0;
}

void tmqSendHbReq(void* param, void* tmrId) {
741 742 743
  int64_t refId = *(int64_t*)param;
  tmq_t*  tmq = taosAcquireRef(tmqMgmt.rsetId, refId);
  if (tmq == NULL) {
L
Liu Jicong 已提交
744
    taosMemoryFree(param);
745 746
    return;
  }
D
dapan1121 已提交
747 748 749 750 751

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

L
Liu Jicong 已提交
752
  int32_t tlen = tSerializeSMqHbReq(NULL, 0, &req);
D
dapan1121 已提交
753 754 755 756
  if (tlen < 0) {
    tscError("tSerializeSMqHbReq failed");
    return;
  }
L
Liu Jicong 已提交
757
  void* pReq = taosMemoryCalloc(1, tlen);
D
dapan1121 已提交
758 759 760 761 762 763 764 765 766
  if (tlen < 0) {
    tscError("failed to malloc MqHbReq msg, size:%d", tlen);
    return;
  }
  if (tSerializeSMqHbReq(pReq, tlen, &req) < 0) {
    tscError("tSerializeSMqHbReq %d failed", tlen);
    taosMemoryFree(pReq);
    return;
  }
767 768 769 770

  SMsgSendInfo* sendInfo = taosMemoryCalloc(1, sizeof(SMsgSendInfo));
  if (sendInfo == NULL) {
    taosMemoryFree(pReq);
L
Liu Jicong 已提交
771
    goto OVER;
772 773 774
  }
  sendInfo->msgInfo = (SDataBuf){
      .pData = pReq,
D
dapan1121 已提交
775
      .len = tlen,
776 777 778 779 780 781 782
      .handle = NULL,
  };

  sendInfo->requestId = generateRequestId();
  sendInfo->requestObjRefId = 0;
  sendInfo->param = NULL;
  sendInfo->fp = tmqHbCb;
L
Liu Jicong 已提交
783
  sendInfo->msgType = TDMT_MND_TMQ_HB;
784 785 786 787 788 789 790

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

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

OVER:
791
  taosTmrReset(tmqSendHbReq, 1000, param, tmqMgmt.timer, &tmq->hbLiveTimer);
792 793
}

794
int32_t tmqHandleAllDelayedTask(tmq_t* pTmq) {
L
Liu Jicong 已提交
795
  STaosQall* qall = taosAllocateQall();
796
  taosReadAllQitems(pTmq->delayedTask, qall);
L
Liu Jicong 已提交
797

798 799 800 801
  if (qall->numOfItems == 0) {
    taosFreeQall(qall);
    return TSDB_CODE_SUCCESS;
  }
802

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

807
  while (pTaskType != NULL) {
808
    if (*pTaskType == TMQ_DELAYED_TASK__ASK_EP) {
809
      tmqAskEp(pTmq, true);
810 811

      int64_t* pRefId = taosMemoryMalloc(sizeof(int64_t));
812
      *pRefId = pTmq->refId;
813

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

      int64_t* pRefId = taosMemoryMalloc(sizeof(int64_t));
820
      *pRefId = pTmq->refId;
821

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

L
Liu Jicong 已提交
828
    taosFreeQitem(pTaskType);
829
    taosGetQitem(qall, (void**)&pTaskType);
L
Liu Jicong 已提交
830
  }
831

L
Liu Jicong 已提交
832 833 834 835
  taosFreeQall(qall);
  return 0;
}

L
Liu Jicong 已提交
836 837 838 839 840 841 842 843 844 845 846 847 848 849 850 851 852 853 854 855 856 857 858 859 860 861 862
static void tmqFreeRspWrapper(SMqRspWrapper* rspWrapper) {
  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;
    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;
    taosMemoryFree(pRsp->metaRsp.metaRsp);
  } else if (rspWrapper->tmqRspType == TMQ_MSG_TYPE__TAOSX_RSP) {
    SMqPollRspWrapper* pRsp = (SMqPollRspWrapper*)rspWrapper;
    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);
  }
}

L
Liu Jicong 已提交
863
void tmqClearUnhandleMsg(tmq_t* tmq) {
L
Liu Jicong 已提交
864
  SMqRspWrapper* rspWrapper = NULL;
L
Liu Jicong 已提交
865
  while (1) {
L
Liu Jicong 已提交
866 867 868 869 870
    taosGetQitem(tmq->qall, (void**)&rspWrapper);
    if (rspWrapper) {
      tmqFreeRspWrapper(rspWrapper);
      taosFreeQitem(rspWrapper);
    } else {
L
Liu Jicong 已提交
871
      break;
L
Liu Jicong 已提交
872
    }
L
Liu Jicong 已提交
873 874
  }

L
Liu Jicong 已提交
875
  rspWrapper = NULL;
L
Liu Jicong 已提交
876 877
  taosReadAllQitems(tmq->mqueue, tmq->qall);
  while (1) {
L
Liu Jicong 已提交
878 879 880 881 882
    taosGetQitem(tmq->qall, (void**)&rspWrapper);
    if (rspWrapper) {
      tmqFreeRspWrapper(rspWrapper);
      taosFreeQitem(rspWrapper);
    } else {
L
Liu Jicong 已提交
883
      break;
L
Liu Jicong 已提交
884
    }
L
Liu Jicong 已提交
885 886 887
  }
}

D
dapan1121 已提交
888
int32_t tmqSubscribeCb(void* param, SDataBuf* pMsg, int32_t code) {
L
Liu Jicong 已提交
889 890
  SMqSubscribeCbParam* pParam = (SMqSubscribeCbParam*)param;
  pParam->rspErr = code;
dengyihao's avatar
dengyihao 已提交
891 892

  taosMemoryFree(pMsg->pEpSet);
L
Liu Jicong 已提交
893 894 895
  tsem_post(&pParam->rspSem);
  return 0;
}
896

L
Liu Jicong 已提交
897
int32_t tmq_subscription(tmq_t* tmq, tmq_list_t** topics) {
X
Xiaoyu Wang 已提交
898 899 900 901
  if (*topics == NULL) {
    *topics = tmq_list_new();
  }
  for (int i = 0; i < taosArrayGetSize(tmq->clientTopics); i++) {
L
Liu Jicong 已提交
902
    SMqClientTopic* topic = taosArrayGet(tmq->clientTopics, i);
L
Liu Jicong 已提交
903
    tmq_list_append(*topics, strchr(topic->topicName, '.') + 1);
X
Xiaoyu Wang 已提交
904
  }
L
Liu Jicong 已提交
905
  return 0;
X
Xiaoyu Wang 已提交
906 907
}

L
Liu Jicong 已提交
908
int32_t tmq_unsubscribe(tmq_t* tmq) {
L
Liu Jicong 已提交
909 910
  int32_t     rsp;
  int32_t     retryCnt = 0;
L
Liu Jicong 已提交
911
  tmq_list_t* lst = tmq_list_new();
L
Liu Jicong 已提交
912 913 914 915 916 917 918 919 920 921
  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 已提交
922 923
  tmq_list_destroy(lst);
  return rsp;
X
Xiaoyu Wang 已提交
924 925
}

926 927
void tmqFreeImpl(void* handle) {
  tmq_t* tmq = (tmq_t*)handle;
L
Liu Jicong 已提交
928

929
  // TODO stop timer
L
Liu Jicong 已提交
930 931 932 933
  if (tmq->mqueue) {
    tmqClearUnhandleMsg(tmq);
    taosCloseQueue(tmq->mqueue);
  }
L
Liu Jicong 已提交
934

H
Haojun Liao 已提交
935 936 937 938 939
  if (tmq->delayedTask) {
    taosCloseQueue(tmq->delayedTask);
  }

  taosFreeQall(tmq->qall);
940
  tsem_destroy(&tmq->rspSem);
H
Haojun Liao 已提交
941
  taosThreadMutexDestroy(&tmq->lock);
L
Liu Jicong 已提交
942

943 944 945
  int32_t sz = taosArrayGetSize(tmq->clientTopics);
  for (int32_t i = 0; i < sz; i++) {
    SMqClientTopic* pTopic = taosArrayGet(tmq->clientTopics, i);
L
Liu Jicong 已提交
946
    taosMemoryFreeClear(pTopic->schema.pSchema);
947 948
    taosArrayDestroy(pTopic->vgs);
  }
H
Haojun Liao 已提交
949

950 951 952
  taosArrayDestroy(tmq->clientTopics);
  taos_close_internal(tmq->pTscObj);
  taosMemoryFree(tmq);
L
Liu Jicong 已提交
953 954
}

955 956 957 958 959 960 961 962 963
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);
964
  if (tmqMgmt.rsetId < 0) {
965 966 967 968
    tmqInitRes = terrno;
  }
}

L
Liu Jicong 已提交
969
tmq_t* tmq_consumer_new(tmq_conf_t* conf, char* errstr, int32_t errstrLen) {
970 971 972 973
  taosThreadOnce(&tmqInit, tmqMgmtInit);
  if (tmqInitRes != 0) {
    terrno = tmqInitRes;
    return NULL;
L
Liu Jicong 已提交
974 975
  }

L
Liu Jicong 已提交
976 977
  tmq_t* pTmq = taosMemoryCalloc(1, sizeof(tmq_t));
  if (pTmq == NULL) {
L
Liu Jicong 已提交
978
    terrno = TSDB_CODE_OUT_OF_MEMORY;
979
    tscError("failed to create consumer, groupId:%s, code:%s", conf->groupId, terrstr());
L
Liu Jicong 已提交
980 981
    return NULL;
  }
L
Liu Jicong 已提交
982

L
Liu Jicong 已提交
983 984 985
  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 已提交
986 987 988
  pTmq->clientTopics = taosArrayInit(0, sizeof(SMqClientTopic));
  pTmq->mqueue = taosOpenQueue();
  pTmq->delayedTask = taosOpenQueue();
H
Haojun Liao 已提交
989
  pTmq->qall = taosAllocateQall();
L
Liu Jicong 已提交
990

H
Haojun Liao 已提交
991
  taosThreadMutexInit(&pTmq->lock, NULL);
X
Xiaoyu Wang 已提交
992 993
  if (pTmq->clientTopics == NULL || pTmq->mqueue == NULL || pTmq->qall == NULL || pTmq->delayedTask == NULL ||
      conf->groupId[0] == 0) {
L
Liu Jicong 已提交
994
    terrno = TSDB_CODE_OUT_OF_MEMORY;
995
    tscError("consumer:0x%" PRIx64 " setup failed since %s, groupId:%s", pTmq->consumerId, terrstr(),
S
Shengliang Guan 已提交
996
             pTmq->groupId);
997
    goto _failed;
L
Liu Jicong 已提交
998
  }
L
Liu Jicong 已提交
999

L
Liu Jicong 已提交
1000 1001
  // init status
  pTmq->status = TMQ_CONSUMER_STATUS__INIT;
L
Liu Jicong 已提交
1002 1003
  pTmq->pollCnt = 0;
  pTmq->epoch = 0;
L
Liu Jicong 已提交
1004 1005
  /*pTmq->epStatus = 0;*/
  /*pTmq->epSkipCnt = 0;*/
L
Liu Jicong 已提交
1006

L
Liu Jicong 已提交
1007 1008 1009
  // set conf
  strcpy(pTmq->clientId, conf->clientId);
  strcpy(pTmq->groupId, conf->groupId);
1010
  pTmq->withTbName = conf->withTbName;
L
Liu Jicong 已提交
1011
  pTmq->useSnapshot = conf->snapEnable;
L
Liu Jicong 已提交
1012
  pTmq->autoCommit = conf->autoCommit;
L
Liu Jicong 已提交
1013
  pTmq->autoCommitInterval = conf->autoCommitInterval;
L
Liu Jicong 已提交
1014 1015
  pTmq->commitCb = conf->commitCb;
  pTmq->commitCbUserParam = conf->commitCbUserParam;
L
Liu Jicong 已提交
1016 1017
  pTmq->resetOffsetCfg = conf->resetOffset;

1018 1019
  pTmq->hbBgEnable = conf->hbBgEnable;

L
Liu Jicong 已提交
1020
  // assign consumerId
L
Liu Jicong 已提交
1021
  pTmq->consumerId = tGenIdPI64();
X
Xiaoyu Wang 已提交
1022

L
Liu Jicong 已提交
1023 1024
  // init semaphore
  if (tsem_init(&pTmq->rspSem, 0, 0) != 0) {
1025
    tscError("consumer:0x %" PRIx64 " setup failed since %s, consumer group %s", pTmq->consumerId, terrstr(),
S
Shengliang Guan 已提交
1026
             pTmq->groupId);
1027
    goto _failed;
L
Liu Jicong 已提交
1028
  }
L
Liu Jicong 已提交
1029

L
Liu Jicong 已提交
1030 1031 1032
  // init connection
  pTmq->pTscObj = taos_connect_internal(conf->ip, user, pass, NULL, NULL, conf->port, CONN_TYPE__TMQ);
  if (pTmq->pTscObj == NULL) {
1033
    tscError("consumer:0x%" PRIx64 " setup failed since %s, groupId:%s", pTmq->consumerId, terrstr(), pTmq->groupId);
L
Liu Jicong 已提交
1034
    tsem_destroy(&pTmq->rspSem);
1035
    goto _failed;
L
Liu Jicong 已提交
1036
  }
L
Liu Jicong 已提交
1037

1038 1039
  pTmq->refId = taosAddRef(tmqMgmt.rsetId, pTmq);
  if (pTmq->refId < 0) {
1040
    goto _failed;
1041 1042
  }

1043
  if (pTmq->hbBgEnable) {
L
Liu Jicong 已提交
1044 1045
    int64_t* pRefId = taosMemoryMalloc(sizeof(int64_t));
    *pRefId = pTmq->refId;
1046
    pTmq->hbLiveTimer = taosTmrStart(tmqSendHbReq, 1000, pRefId, tmqMgmt.timer);
1047 1048
  }

1049 1050 1051 1052 1053 1054
  char         buf[80] = {0};
  STqOffsetVal offset = {.type = pTmq->resetOffsetCfg};
  tFormatOffset(buf, tListLen(buf), &offset);
  tscInfo("consumer:0x%" PRIx64 " is setup, groupId:%s, snapshot:%d, autoCommit:%d, commitInterval:%dms, offset:%s, backgroudHB:%d",
          pTmq->consumerId, pTmq->groupId, pTmq->useSnapshot, pTmq->autoCommit, pTmq->autoCommitInterval, buf,
          pTmq->hbBgEnable);
L
Liu Jicong 已提交
1055

1056
  return pTmq;
1057

1058 1059
_failed:
  tmqFreeImpl(pTmq);
L
Liu Jicong 已提交
1060
  return NULL;
1061 1062
}

L
Liu Jicong 已提交
1063
int32_t tmq_subscribe(tmq_t* tmq, const tmq_list_t* topic_list) {
L
Liu Jicong 已提交
1064 1065 1066
  const SArray*   container = &topic_list->container;
  int32_t         sz = taosArrayGetSize(container);
  void*           buf = NULL;
L
Liu Jicong 已提交
1067
  SMsgSendInfo*   sendInfo = NULL;
L
Liu Jicong 已提交
1068
  SCMSubscribeReq req = {0};
1069
  int32_t         code = 0;
1070

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

1073
  req.consumerId = tmq->consumerId;
L
Liu Jicong 已提交
1074
  tstrncpy(req.clientId, tmq->clientId, 256);
L
Liu Jicong 已提交
1075
  tstrncpy(req.cgroup, tmq->groupId, TSDB_CGROUP_LEN);
1076 1077
  req.topicNames = taosArrayInit(sz, sizeof(void*));

1078 1079 1080 1081
  if (req.topicNames == NULL) {
    code = TSDB_CODE_OUT_OF_MEMORY;
    goto FAIL;
  }
1082

L
Liu Jicong 已提交
1083 1084
  for (int32_t i = 0; i < sz; i++) {
    char* topic = taosArrayGetP(container, i);
1085 1086

    SName name = {0};
L
Liu Jicong 已提交
1087 1088 1089 1090
    tNameSetDbName(&name, tmq->pTscObj->acctId, topic, strlen(topic));
    char* topicFName = taosMemoryCalloc(1, TSDB_TOPIC_FNAME_LEN);
    if (topicFName == NULL) {
      goto FAIL;
1091 1092
    }

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

    taosArrayPush(req.topicNames, &topicFName);
1097 1098
  }

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

L
Liu Jicong 已提交
1101
  buf = taosMemoryMalloc(tlen);
1102 1103 1104 1105
  if (buf == NULL) {
    code = TSDB_CODE_OUT_OF_MEMORY;
    goto FAIL;
  }
L
Liu Jicong 已提交
1106

1107 1108 1109
  void* abuf = buf;
  tSerializeSCMSubscribeReq(&abuf, &req);

L
Liu Jicong 已提交
1110
  sendInfo = taosMemoryCalloc(1, sizeof(SMsgSendInfo));
1111 1112 1113 1114
  if (sendInfo == NULL) {
    code = TSDB_CODE_OUT_OF_MEMORY;
    goto FAIL;
  }
1115

X
Xiaoyu Wang 已提交
1116
  SMqSubscribeCbParam param = {
L
Liu Jicong 已提交
1117
      .rspErr = 0,
1118 1119
      .refId = tmq->refId,
      .epoch = tmq->epoch,
X
Xiaoyu Wang 已提交
1120
  };
L
Liu Jicong 已提交
1121

1122 1123 1124
  if (tsem_init(&param.rspSem, 0, 0) != 0) {
    goto FAIL;
  }
L
Liu Jicong 已提交
1125 1126

  sendInfo->msgInfo = (SDataBuf){
X
Xiaoyu Wang 已提交
1127 1128 1129 1130
      .pData = buf,
      .len = tlen,
      .handle = NULL,
  };
1131

L
Liu Jicong 已提交
1132 1133
  sendInfo->requestId = generateRequestId();
  sendInfo->requestObjRefId = 0;
L
Liu Jicong 已提交
1134 1135
  sendInfo->param = &param;
  sendInfo->fp = tmqSubscribeCb;
L
Liu Jicong 已提交
1136
  sendInfo->msgType = TDMT_MND_TMQ_SUBSCRIBE;
L
Liu Jicong 已提交
1137

1138 1139 1140 1141 1142
  SEpSet epSet = getEpSet_s(&tmq->pTscObj->pAppInfo->mgmtEp);

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

L
Liu Jicong 已提交
1143 1144
  // avoid double free if msg is sent
  buf = NULL;
L
Liu Jicong 已提交
1145
  sendInfo = NULL;
L
Liu Jicong 已提交
1146

L
Liu Jicong 已提交
1147 1148
  tsem_wait(&param.rspSem);
  tsem_destroy(&param.rspSem);
1149

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

L
Liu Jicong 已提交
1155
  int32_t retryCnt = 0;
L
Liu Jicong 已提交
1156
  while (TSDB_CODE_MND_CONSUMER_NOT_READY == tmqAskEp(tmq, false)) {
wmmhello's avatar
wmmhello 已提交
1157
    if (retryCnt++ > 40) {
L
Liu Jicong 已提交
1158 1159
      goto FAIL;
    }
1160

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

1165 1166
  // init ep timer
  if (tmq->epTimer == NULL) {
1167 1168 1169
    int64_t* pRefId1 = taosMemoryMalloc(sizeof(int64_t));
    *pRefId1 = tmq->refId;
    tmq->epTimer = taosTmrStart(tmqAssignAskEpTask, 1000, pRefId1, tmqMgmt.timer);
1170
  }
L
Liu Jicong 已提交
1171 1172

  // init auto commit timer
1173
  if (tmq->autoCommit && tmq->commitTimer == NULL) {
1174 1175 1176
    int64_t* pRefId2 = taosMemoryMalloc(sizeof(int64_t));
    *pRefId2 = tmq->refId;
    tmq->commitTimer = taosTmrStart(tmqAssignDelayedCommitTask, tmq->autoCommitInterval, pRefId2, tmqMgmt.timer);
L
Liu Jicong 已提交
1177 1178
  }

L
Liu Jicong 已提交
1179
FAIL:
L
Liu Jicong 已提交
1180
  taosArrayDestroyP(req.topicNames, taosMemoryFree);
L
Liu Jicong 已提交
1181
  taosMemoryFree(buf);
L
Liu Jicong 已提交
1182
  taosMemoryFree(sendInfo);
L
Liu Jicong 已提交
1183

L
Liu Jicong 已提交
1184
  return code;
1185 1186
}

L
Liu Jicong 已提交
1187
void tmq_conf_set_auto_commit_cb(tmq_conf_t* conf, tmq_commit_cb* cb, void* param) {
1188
  conf->commitCb = cb;
L
Liu Jicong 已提交
1189
  conf->commitCbUserParam = param;
L
Liu Jicong 已提交
1190
}
1191

D
dapan1121 已提交
1192
int32_t tmqPollCb(void* param, SDataBuf* pMsg, int32_t code) {
X
Xiaoyu Wang 已提交
1193 1194
  SMqPollCbParam* pParam = (SMqPollCbParam*)param;
  SMqClientVg*    pVg = pParam->pVg;
L
Liu Jicong 已提交
1195
  SMqClientTopic* pTopic = pParam->pTopic;
1196 1197 1198 1199 1200 1201

  tmq_t* tmq = taosAcquireRef(tmqMgmt.rsetId, pParam->refId);
  if (tmq == NULL) {
    tsem_destroy(&pParam->rspSem);
    taosMemoryFree(pParam);
    taosMemoryFree(pMsg->pData);
dengyihao's avatar
dengyihao 已提交
1202
    taosMemoryFree(pMsg->pEpSet);
1203 1204 1205 1206
    terrno = TSDB_CODE_TMQ_CONSUMER_CLOSED;
    return -1;
  }

H
Haojun Liao 已提交
1207 1208 1209 1210
  int32_t  epoch = pParam->epoch;
  int32_t  vgId = pParam->vgId;
  uint64_t requestId = pParam->requestId;

L
Liu Jicong 已提交
1211
  taosMemoryFree(pParam);
H
Haojun Liao 已提交
1212

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

L
Liu Jicong 已提交
1217
    if (pMsg->pData) taosMemoryFree(pMsg->pData);
dengyihao's avatar
dengyihao 已提交
1218 1219
    if (pMsg->pEpSet) taosMemoryFree(pMsg->pEpSet);

H
Haojun Liao 已提交
1220
    // in case of consumer mismatch, wait for 500ms and retry
L
Liu Jicong 已提交
1221
    if (code == TSDB_CODE_TMQ_CONSUMER_MISMATCH) {
1222
      taosMsleep(500);
L
Liu Jicong 已提交
1223
      atomic_store_8(&tmq->status, TMQ_CONSUMER_STATUS__RECOVER);
1224
      tscDebug("consumer:0x%" PRIx64" wait for the re-balance, wait for 500ms and set status to be RECOVER", tmq->consumerId);
H
Haojun Liao 已提交
1225
    } else if (code == TSDB_CODE_TQ_NO_COMMITTED_OFFSET) {
S
Shengliang Guan 已提交
1226
      SMqPollRspWrapper* pRspWrapper = taosAllocateQitem(sizeof(SMqPollRspWrapper), DEF_QITEM, 0);
L
Liu Jicong 已提交
1227
      if (pRspWrapper == NULL) {
H
Haojun Liao 已提交
1228 1229
        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 已提交
1230 1231
        goto CREATE_MSG_FAIL;
      }
H
Haojun Liao 已提交
1232

L
Liu Jicong 已提交
1233 1234 1235 1236 1237 1238
      pRspWrapper->tmqRspType = TMQ_MSG_TYPE__END_RSP;
      /*pRspWrapper->vgHandle = pVg;*/
      /*pRspWrapper->topicHandle = pTopic;*/
      taosWriteQitem(tmq->mqueue, pRspWrapper);
      tsem_post(&tmq->rspSem);
    }
H
Haojun Liao 已提交
1239

L
fix txn  
Liu Jicong 已提交
1240
    goto CREATE_MSG_FAIL;
L
Liu Jicong 已提交
1241 1242
  }

X
Xiaoyu Wang 已提交
1243 1244 1245
  int32_t msgEpoch = ((SMqRspHead*)pMsg->pData)->epoch;
  int32_t tmqEpoch = atomic_load_32(&tmq->epoch);
  if (msgEpoch < tmqEpoch) {
L
Liu Jicong 已提交
1246
    // do not write into queue since updating epoch reset
H
Haojun Liao 已提交
1247 1248 1249
    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);

1250
    tsem_post(&tmq->rspSem);
L
Liu Jicong 已提交
1251
    taosMemoryFree(pMsg->pData);
dengyihao's avatar
dengyihao 已提交
1252
    taosMemoryFree(pMsg->pEpSet);
X
Xiaoyu Wang 已提交
1253 1254 1255 1256
    return 0;
  }

  if (msgEpoch != tmqEpoch) {
H
Haojun Liao 已提交
1257 1258
    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 已提交
1259 1260
  }

L
Liu Jicong 已提交
1261 1262 1263
  // handle meta rsp
  int8_t rspType = ((SMqRspHead*)pMsg->pData)->mqMsgType;

S
Shengliang Guan 已提交
1264
  SMqPollRspWrapper* pRspWrapper = taosAllocateQitem(sizeof(SMqPollRspWrapper), DEF_QITEM, 0);
L
Liu Jicong 已提交
1265
  if (pRspWrapper == NULL) {
L
Liu Jicong 已提交
1266
    taosMemoryFree(pMsg->pData);
dengyihao's avatar
dengyihao 已提交
1267
    taosMemoryFree(pMsg->pEpSet);
H
Haojun Liao 已提交
1268
    tscWarn("consumer:0x%"PRIx64" msg discard from vgId:%d, epoch %d since out of memory", tmq->consumerId, vgId, epoch);
L
fix txn  
Liu Jicong 已提交
1269
    goto CREATE_MSG_FAIL;
L
Liu Jicong 已提交
1270
  }
L
Liu Jicong 已提交
1271

L
Liu Jicong 已提交
1272
  pRspWrapper->tmqRspType = rspType;
L
Liu Jicong 已提交
1273 1274
  pRspWrapper->vgHandle = pVg;
  pRspWrapper->topicHandle = pTopic;
L
Liu Jicong 已提交
1275

L
Liu Jicong 已提交
1276
  if (rspType == TMQ_MSG_TYPE__POLL_RSP) {
L
Liu Jicong 已提交
1277 1278 1279
    SDecoder decoder;
    tDecoderInit(&decoder, POINTER_SHIFT(pMsg->pData, sizeof(SMqRspHead)), pMsg->len - sizeof(SMqRspHead));
    tDecodeSMqDataRsp(&decoder, &pRspWrapper->dataRsp);
wmmhello's avatar
wmmhello 已提交
1280
    tDecoderClear(&decoder);
L
Liu Jicong 已提交
1281
    memcpy(&pRspWrapper->dataRsp, pMsg->pData, sizeof(SMqRspHead));
L
Liu Jicong 已提交
1282

H
Haojun Liao 已提交
1283 1284
    tscDebug("consumer:0x%" PRIx64 " recv poll rsp, vgId:%d, req:%" PRId64 ", rsp:%" PRId64
             " type %d, reqId:0x%" PRIx64,
H
Haojun Liao 已提交
1285
             tmq->consumerId, vgId, pRspWrapper->dataRsp.reqOffset.version, pRspWrapper->dataRsp.rspOffset.version, rspType, requestId);
H
Haojun Liao 已提交
1286

L
Liu Jicong 已提交
1287
  } else if (rspType == TMQ_MSG_TYPE__POLL_META_RSP) {
1288 1289 1290 1291
    SDecoder decoder;
    tDecoderInit(&decoder, POINTER_SHIFT(pMsg->pData, sizeof(SMqRspHead)), pMsg->len - sizeof(SMqRspHead));
    tDecodeSMqMetaRsp(&decoder, &pRspWrapper->metaRsp);
    tDecoderClear(&decoder);
L
Liu Jicong 已提交
1292
    memcpy(&pRspWrapper->metaRsp, pMsg->pData, sizeof(SMqRspHead));
L
Liu Jicong 已提交
1293 1294 1295 1296 1297 1298
  } 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));
L
Liu Jicong 已提交
1299
  }
L
Liu Jicong 已提交
1300

L
Liu Jicong 已提交
1301
  taosMemoryFree(pMsg->pData);
dengyihao's avatar
dengyihao 已提交
1302
  taosMemoryFree(pMsg->pEpSet);
L
Liu Jicong 已提交
1303

H
Haojun Liao 已提交
1304 1305
  tscDebug("consumer:0x%" PRIx64 ", put poll res into mqueue, total in queue:%d, reqId:0x%" PRIx64, tmq->consumerId,
           tmq->mqueue->numOfItems, requestId);
L
Liu Jicong 已提交
1306

L
Liu Jicong 已提交
1307
  taosWriteQitem(tmq->mqueue, pRspWrapper);
1308
  tsem_post(&tmq->rspSem);
L
Liu Jicong 已提交
1309

L
Liu Jicong 已提交
1310
  return 0;
H
Haojun Liao 已提交
1311

L
fix txn  
Liu Jicong 已提交
1312
CREATE_MSG_FAIL:
L
Liu Jicong 已提交
1313
  if (epoch == tmq->epoch) {
L
Liu Jicong 已提交
1314 1315
    atomic_store_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE);
  }
H
Haojun Liao 已提交
1316

1317
  tsem_post(&tmq->rspSem);
L
Liu Jicong 已提交
1318
  return -1;
1319 1320
}

H
Haojun Liao 已提交
1321 1322 1323 1324 1325 1326 1327 1328 1329 1330 1331 1332 1333 1334 1335 1336 1337
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 已提交
1338 1339

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

H
Haojun Liao 已提交
1342 1343 1344 1345 1346 1347 1348 1349 1350 1351 1352 1353 1354 1355 1356 1357 1358 1359 1360 1361 1362 1363 1364 1365 1366 1367 1368 1369
    STqOffsetVal  offsetNew = {.type = tmq->resetOffsetCfg};
    if (pOffset != NULL) {
      offsetNew = *pOffset;
    }

    SMqClientVg clientVg = {
        .pollCnt = 0,
        .currentOffset = offsetNew,
        .vgId = pVgEp->vgId,
        .epSet = pVgEp->epSet,
        .vgStatus = TMQ_VG_STATUS__IDLE,
        .vgSkipCnt = 0,
    };

    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) {
1370 1371
  bool set = false;

1372
  int32_t topicNumCur = taosArrayGetSize(tmq->clientTopics);
1373
  int32_t topicNumGet = taosArrayGetSize(pRsp->topics);
1374

X
Xiaoyu Wang 已提交
1375 1376
  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",
1377
           tmq->consumerId, tmq->epoch, epoch, topicNumGet, topicNumCur);
1378 1379 1380 1381 1382 1383

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

H
Haojun Liao 已提交
1384 1385
  SHashObj* pVgOffsetHashMap = taosHashInit(64, MurmurHash3_32, false, HASH_NO_LOCK);
  if (pVgOffsetHashMap == NULL) {
1386 1387 1388
    taosArrayDestroy(newTopics);
    return false;
  }
1389

H
Haojun Liao 已提交
1390
  // todo extract method
1391 1392 1393 1394 1395
  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);
1396
      tscDebug("consumer:0x%" PRIx64 ", new vg num: %d", tmq->consumerId, vgNumCur);
1397 1398
      for (int32_t j = 0; j < vgNumCur; j++) {
        SMqClientVg* pVgCur = taosArrayGet(pTopicCur->vgs, j);
H
Haojun Liao 已提交
1399 1400
        makeTopicVgroupKey(vgKey, pTopicCur->topicName, pVgCur->vgId);

L
Liu Jicong 已提交
1401
        char buf[80];
L
Liu Jicong 已提交
1402
        tFormatOffset(buf, 80, &pVgCur->currentOffset);
H
Haojun Liao 已提交
1403
        tscDebug("consumer:0x%" PRIx64 ", epoch:%d vgId:%d vgKey:%s, offset:%s", tmq->consumerId, epoch,
L
Liu Jicong 已提交
1404
                 pVgCur->vgId, vgKey, buf);
H
Haojun Liao 已提交
1405
        taosHashPut(pVgOffsetHashMap, vgKey, strlen(vgKey), &pVgCur->currentOffset, sizeof(STqOffsetVal));
1406 1407 1408 1409 1410 1411 1412
      }
    }
  }

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

H
Haojun Liao 已提交
1417 1418 1419
  taosHashCleanup(pVgOffsetHashMap);

  taosThreadMutexLock(&tmq->lock);
1420
  // destroy current buffered existed topics info
1421
  if (tmq->clientTopics) {
H
Haojun Liao 已提交
1422
    taosArrayDestroyEx(tmq->clientTopics, freeClientVgInfo);
X
Xiaoyu Wang 已提交
1423
  }
1424

H
Haojun Liao 已提交
1425 1426
  tmq->clientTopics = newTopics;
  taosThreadMutexUnlock(&tmq->lock);
1427

H
Haojun Liao 已提交
1428 1429
  int8_t flag = (topicNumGet == 0)? TMQ_CONSUMER_STATUS__NO_TOPIC:TMQ_CONSUMER_STATUS__READY;
  atomic_store_8(&tmq->status, flag);
X
Xiaoyu Wang 已提交
1430
  atomic_store_32(&tmq->epoch, epoch);
H
Haojun Liao 已提交
1431

1432
  tscDebug("consumer:0x%" PRIx64 " update topic info completed", tmq->consumerId);
X
Xiaoyu Wang 已提交
1433 1434 1435
  return set;
}

H
Haojun Liao 已提交
1436
static int32_t tmqAskEpCb(void* param, SDataBuf* pMsg, int32_t code) {
1437
  SMqAskEpCbParam* pParam = (SMqAskEpCbParam*)param;
L
Liu Jicong 已提交
1438
  int8_t           async = pParam->async;
1439 1440 1441 1442 1443 1444 1445 1446 1447
  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 已提交
1448
    taosMemoryFree(pMsg->pEpSet);
1449 1450 1451 1452
    terrno = TSDB_CODE_TMQ_CONSUMER_CLOSED;
    return -1;
  }

L
Liu Jicong 已提交
1453
  pParam->code = code;
H
Haojun Liao 已提交
1454
  if (code != TSDB_CODE_SUCCESS) {
X
Xiaoyu Wang 已提交
1455 1456
    tscError("consumer:0x%" PRIx64 ", get topic endpoint error, async:%d, code:%s", tmq->consumerId, pParam->async,
             tstrerror(code));
L
Liu Jicong 已提交
1457
    goto END;
1458
  }
L
Liu Jicong 已提交
1459

L
Liu Jicong 已提交
1460
  // tmq's epoch is monotonically increase,
L
Liu Jicong 已提交
1461
  // so it's safe to discard any old epoch msg.
L
Liu Jicong 已提交
1462
  // Epoch will only increase when received newer epoch ep msg
L
Liu Jicong 已提交
1463 1464 1465
  SMqRspHead* head = pMsg->pData;
  int32_t     epoch = atomic_load_32(&tmq->epoch);
  if (head->epoch <= epoch) {
1466 1467
    tscDebug("consumer:0x%" PRIx64 ", recv ep, msg epoch %d, current epoch %d, no need to update local ep",
             tmq->consumerId, head->epoch, epoch);
L
Liu Jicong 已提交
1468
    goto END;
1469
  }
L
Liu Jicong 已提交
1470

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

L
Liu Jicong 已提交
1474
  if (!async) {
L
Liu Jicong 已提交
1475 1476
    SMqAskEpRsp rsp;
    tDecodeSMqAskEpRsp(POINTER_SHIFT(pMsg->pData, sizeof(SMqRspHead)), &rsp);
L
Liu Jicong 已提交
1477
    tmqUpdateEp(tmq, head->epoch, &rsp);
L
Liu Jicong 已提交
1478
    tDeleteSMqAskEpRsp(&rsp);
X
Xiaoyu Wang 已提交
1479
  } else {
S
Shengliang Guan 已提交
1480
    SMqAskEpRspWrapper* pWrapper = taosAllocateQitem(sizeof(SMqAskEpRspWrapper), DEF_QITEM, 0);
L
Liu Jicong 已提交
1481
    if (pWrapper == NULL) {
X
Xiaoyu Wang 已提交
1482
      terrno = TSDB_CODE_OUT_OF_MEMORY;
L
Liu Jicong 已提交
1483 1484
      code = -1;
      goto END;
X
Xiaoyu Wang 已提交
1485
    }
1486

L
Liu Jicong 已提交
1487 1488 1489
    pWrapper->tmqRspType = TMQ_MSG_TYPE__EP_RSP;
    pWrapper->epoch = head->epoch;
    memcpy(&pWrapper->msg, pMsg->pData, sizeof(SMqRspHead));
L
Liu Jicong 已提交
1490
    tDecodeSMqAskEpRsp(POINTER_SHIFT(pMsg->pData, sizeof(SMqRspHead)), &pWrapper->msg);
L
Liu Jicong 已提交
1491

L
Liu Jicong 已提交
1492
    taosWriteQitem(tmq->mqueue, pWrapper);
1493
    tsem_post(&tmq->rspSem);
1494
  }
L
Liu Jicong 已提交
1495 1496

END:
L
Liu Jicong 已提交
1497
  if (!async) {
L
Liu Jicong 已提交
1498
    tsem_post(&pParam->rspSem);
L
Liu Jicong 已提交
1499 1500
  } else {
    taosMemoryFree(pParam);
L
Liu Jicong 已提交
1501
  }
dengyihao's avatar
dengyihao 已提交
1502 1503

  taosMemoryFree(pMsg->pEpSet);
L
Liu Jicong 已提交
1504
  taosMemoryFree(pMsg->pData);
L
Liu Jicong 已提交
1505
  return code;
1506 1507
}

L
Liu Jicong 已提交
1508
void tmqBuildConsumeReqImpl(SMqPollReq* pReq, tmq_t* tmq, int64_t timeout, SMqClientTopic* pTopic, SMqClientVg* pVg) {
L
Liu Jicong 已提交
1509 1510 1511 1512
  int32_t groupLen = strlen(tmq->groupId);
  memcpy(pReq->subKey, tmq->groupId, groupLen);
  pReq->subKey[groupLen] = TMQ_SEPARATOR;
  strcpy(pReq->subKey + groupLen + 1, pTopic->topicName);
1513

1514
  pReq->withTbName = tmq->withTbName;
L
Liu Jicong 已提交
1515
  pReq->consumerId = tmq->consumerId;
1516
  pReq->timeout = timeout;
X
Xiaoyu Wang 已提交
1517
  pReq->epoch = tmq->epoch;
L
Liu Jicong 已提交
1518
  /*pReq->currentOffset = reqOffset;*/
L
Liu Jicong 已提交
1519
  pReq->reqOffset = pVg->currentOffset;
D
dapan1121 已提交
1520
  pReq->head.vgId = pVg->vgId;
1521 1522
  pReq->useSnapshot = tmq->useSnapshot;
  pReq->reqId = generateRequestId();
1523 1524
}

L
Liu Jicong 已提交
1525 1526
SMqMetaRspObj* tmqBuildMetaRspFromWrapper(SMqPollRspWrapper* pWrapper) {
  SMqMetaRspObj* pRspObj = taosMemoryCalloc(1, sizeof(SMqMetaRspObj));
L
Liu Jicong 已提交
1527
  pRspObj->resType = RES_TYPE__TMQ_META;
L
Liu Jicong 已提交
1528 1529 1530 1531 1532 1533 1534 1535
  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;
}

L
Liu Jicong 已提交
1536 1537 1538
SMqRspObj* tmqBuildRspFromWrapper(SMqPollRspWrapper* pWrapper) {
  SMqRspObj* pRspObj = taosMemoryCalloc(1, sizeof(SMqRspObj));
  pRspObj->resType = RES_TYPE__TMQ;
L
Liu Jicong 已提交
1539 1540
  tstrncpy(pRspObj->topic, pWrapper->topicHandle->topicName, TSDB_TOPIC_FNAME_LEN);
  tstrncpy(pRspObj->db, pWrapper->topicHandle->db, TSDB_DB_FNAME_LEN);
L
Liu Jicong 已提交
1541
  pRspObj->vgId = pWrapper->vgHandle->vgId;
L
Liu Jicong 已提交
1542
  pRspObj->resIter = -1;
L
Liu Jicong 已提交
1543
  memcpy(&pRspObj->rsp, &pWrapper->dataRsp, sizeof(SMqDataRsp));
L
Liu Jicong 已提交
1544

L
Liu Jicong 已提交
1545 1546
  pRspObj->resInfo.totalRows = 0;
  pRspObj->resInfo.precision = TSDB_TIME_PRECISION_MILLI;
L
Liu Jicong 已提交
1547
  if (!pWrapper->dataRsp.withSchema) {
L
Liu Jicong 已提交
1548 1549
    setResSchemaInfo(&pRspObj->resInfo, pWrapper->topicHandle->schema.pSchema, pWrapper->topicHandle->schema.nCols);
  }
L
Liu Jicong 已提交
1550

L
Liu Jicong 已提交
1551
  return pRspObj;
X
Xiaoyu Wang 已提交
1552 1553
}

L
Liu Jicong 已提交
1554 1555
SMqTaosxRspObj* tmqBuildTaosxRspFromWrapper(SMqPollRspWrapper* pWrapper) {
  SMqTaosxRspObj* pRspObj = taosMemoryCalloc(1, sizeof(SMqTaosxRspObj));
1556
  pRspObj->resType = RES_TYPE__TMQ_METADATA;
L
Liu Jicong 已提交
1557 1558 1559 1560
  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;
1561
  memcpy(&pRspObj->rsp, &pWrapper->taosxRsp, sizeof(STaosxRsp));
L
Liu Jicong 已提交
1562 1563 1564

  pRspObj->resInfo.totalRows = 0;
  pRspObj->resInfo.precision = TSDB_TIME_PRECISION_MILLI;
1565
  if (!pWrapper->taosxRsp.withSchema) {
L
Liu Jicong 已提交
1566 1567 1568 1569 1570 1571
    setResSchemaInfo(&pRspObj->resInfo, pWrapper->topicHandle->schema.pSchema, pWrapper->topicHandle->schema.nCols);
  }

  return pRspObj;
}

1572 1573 1574 1575 1576 1577 1578 1579 1580 1581 1582 1583 1584 1585 1586 1587 1588 1589 1590 1591 1592 1593 1594 1595 1596 1597 1598 1599 1600 1601 1602 1603 1604
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 已提交
1605
  pParam->pVg = pVg;  // pVg may be released,fix it
1606 1607
  pParam->pTopic = pTopic;
  pParam->vgId = pVg->vgId;
H
Haojun Liao 已提交
1608
  pParam->requestId = req.reqId;
1609 1610 1611 1612 1613 1614 1615 1616 1617 1618 1619 1620 1621 1622 1623 1624 1625 1626 1627 1628 1629 1630 1631 1632

  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 已提交
1633
  tscDebug("consumer:0x%" PRIx64 " send poll to %s vgId:%d, epoch %d, req:%s, reqId:0x%" PRIx64,
1634 1635 1636 1637 1638 1639 1640 1641 1642
           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;
}

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

  for (int i = 0; i < numOfTopics; i++) {
X
Xiaoyu Wang 已提交
1649
    SMqClientTopic* pTopic = taosArrayGet(tmq->clientTopics, i);
X
Xiaoyu Wang 已提交
1650
    int32_t         numOfVg = taosArrayGetSize(pTopic->vgs);
1651 1652

    for (int j = 0; j < numOfVg; j++) {
X
Xiaoyu Wang 已提交
1653 1654
      SMqClientVg* pVg = taosArrayGet(pTopic->vgs, j);
      int32_t      vgStatus = atomic_val_compare_exchange_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE, TMQ_VG_STATUS__WAIT);
1655
      if (vgStatus == TMQ_VG_STATUS__WAIT) {
L
Liu Jicong 已提交
1656
        int32_t vgSkipCnt = atomic_add_fetch_32(&pVg->vgSkipCnt, 1);
H
Haojun Liao 已提交
1657
        tscTrace("consumer:0x%" PRIx64 " epoch %d wait poll-rsp, skip vgId:%d skip cnt %d", tmq->consumerId, tmq->epoch,
X
Xiaoyu Wang 已提交
1658
                 pVg->vgId, vgSkipCnt);
X
Xiaoyu Wang 已提交
1659
        continue;
L
temp  
Liu Jicong 已提交
1660 1661 1662 1663
#if 0
        if (skipCnt < 30000) {
          continue;
        } else {
1664
        tscDebug("consumer:0x%" PRIx64 ",skip vgId:%d skip too much reset", tmq->consumerId, pVg->vgId);
L
temp  
Liu Jicong 已提交
1665 1666
        }
#endif
X
Xiaoyu Wang 已提交
1667
      }
1668

L
Liu Jicong 已提交
1669
      atomic_store_32(&pVg->vgSkipCnt, 0);
1670 1671 1672
      int32_t code = doTmqPollImpl(tmq, pTopic, pVg, timeout);
      if (code != TSDB_CODE_SUCCESS) {
        return code;
D
dapan1121 已提交
1673
      }
X
Xiaoyu Wang 已提交
1674 1675
    }
  }
1676

X
Xiaoyu Wang 已提交
1677 1678 1679
  return 0;
}

H
Haojun Liao 已提交
1680
static int32_t tmqHandleNoPollRsp(tmq_t* tmq, SMqRspWrapper* rspWrapper, bool* pReset) {
L
Liu Jicong 已提交
1681
  if (rspWrapper->tmqRspType == TMQ_MSG_TYPE__EP_RSP) {
L
fix  
Liu Jicong 已提交
1682
    /*printf("ep %d %d\n", rspMsg->head.epoch, tmq->epoch);*/
L
Liu Jicong 已提交
1683 1684
    if (rspWrapper->epoch > atomic_load_32(&tmq->epoch)) {
      SMqAskEpRspWrapper* pEpRspWrapper = (SMqAskEpRspWrapper*)rspWrapper;
L
Liu Jicong 已提交
1685
      SMqAskEpRsp*        rspMsg = &pEpRspWrapper->msg;
L
Liu Jicong 已提交
1686
      tmqUpdateEp(tmq, rspWrapper->epoch, rspMsg);
L
temp  
Liu Jicong 已提交
1687
      /*tmqClearUnhandleMsg(tmq);*/
L
Liu Jicong 已提交
1688
      tDeleteSMqAskEpRsp(rspMsg);
X
Xiaoyu Wang 已提交
1689 1690
      *pReset = true;
    } else {
L
Liu Jicong 已提交
1691
      tmqFreeRspWrapper(rspWrapper);
X
Xiaoyu Wang 已提交
1692 1693 1694 1695 1696 1697 1698 1699
      *pReset = false;
    }
  } else {
    return -1;
  }
  return 0;
}

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

X
Xiaoyu Wang 已提交
1703
  while (1) {
L
Liu Jicong 已提交
1704 1705
    SMqRspWrapper* rspWrapper = NULL;
    taosGetQitem(tmq->qall, (void**)&rspWrapper);
1706

L
Liu Jicong 已提交
1707
    if (rspWrapper == NULL) {
L
Liu Jicong 已提交
1708
      taosReadAllQitems(tmq->mqueue, tmq->qall);
L
Liu Jicong 已提交
1709
      taosGetQitem(tmq->qall, (void**)&rspWrapper);
L
Liu Jicong 已提交
1710 1711 1712 1713

      if (rspWrapper == NULL) {
        return NULL;
      }
X
Xiaoyu Wang 已提交
1714 1715
    }

L
Liu Jicong 已提交
1716 1717 1718
    if (rspWrapper->tmqRspType == TMQ_MSG_TYPE__END_RSP) {
      taosFreeQitem(rspWrapper);
      terrno = TSDB_CODE_TQ_NO_COMMITTED_OFFSET;
H
Haojun Liao 已提交
1719
      tscError("consumer:0x%" PRIx64 " unexpected rsp from poll, code:%s", tmq->consumerId, tstrerror(terrno));
L
Liu Jicong 已提交
1720 1721
      return NULL;
    } else if (rspWrapper->tmqRspType == TMQ_MSG_TYPE__POLL_RSP) {
L
Liu Jicong 已提交
1722
      SMqPollRspWrapper* pollRspWrapper = (SMqPollRspWrapper*)rspWrapper;
H
Haojun Liao 已提交
1723

L
Liu Jicong 已提交
1724
      /*atomic_sub_fetch_32(&tmq->readyRequest, 1);*/
1725
      int32_t consumerEpoch = atomic_load_32(&tmq->epoch);
L
Liu Jicong 已提交
1726
      if (pollRspWrapper->dataRsp.head.epoch == consumerEpoch) {
L
Liu Jicong 已提交
1727
        SMqClientVg* pVg = pollRspWrapper->vgHandle;
L
Liu Jicong 已提交
1728
        pVg->currentOffset = pollRspWrapper->dataRsp.rspOffset;
X
Xiaoyu Wang 已提交
1729
        atomic_store_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE);
H
Haojun Liao 已提交
1730

L
Liu Jicong 已提交
1731
        if (pollRspWrapper->dataRsp.blockNum == 0) {
H
Haojun Liao 已提交
1732
          tscDebug("consumer:0x%" PRIx64 " empty block received, vgId:%d", tmq->consumerId, pVg->vgId);
L
Liu Jicong 已提交
1733 1734
          taosFreeQitem(pollRspWrapper);
          rspWrapper = NULL;
L
temp  
Liu Jicong 已提交
1735 1736
          continue;
        }
H
Haojun Liao 已提交
1737

L
Liu Jicong 已提交
1738
        // build rsp
H
Haojun Liao 已提交
1739 1740
        char buf[80];
        tFormatOffset(buf, 80, &pVg->currentOffset);
L
Liu Jicong 已提交
1741
        SMqRspObj* pRsp = tmqBuildRspFromWrapper(pollRspWrapper);
H
Haojun Liao 已提交
1742 1743 1744
        tscDebug("consumer:0x%" PRIx64 " process poll rsp, vgId:%d, offset:%s, blocks:%d", tmq->consumerId,
                     pVg->vgId, buf, pollRspWrapper->dataRsp.blockNum);

L
Liu Jicong 已提交
1745
        taosFreeQitem(pollRspWrapper);
L
Liu Jicong 已提交
1746
        return pRsp;
X
Xiaoyu Wang 已提交
1747
      } else {
X
Xiaoyu Wang 已提交
1748
        tscDebug("consumer:0x%" PRIx64 " msg discard since epoch mismatch: msg epoch %d, consumer epoch %d",
1749
                 tmq->consumerId, pollRspWrapper->dataRsp.head.epoch, consumerEpoch);
L
Liu Jicong 已提交
1750
        tmqFreeRspWrapper(rspWrapper);
L
Liu Jicong 已提交
1751 1752 1753 1754 1755
        taosFreeQitem(pollRspWrapper);
      }
    } else if (rspWrapper->tmqRspType == TMQ_MSG_TYPE__POLL_META_RSP) {
      SMqPollRspWrapper* pollRspWrapper = (SMqPollRspWrapper*)rspWrapper;
      int32_t            consumerEpoch = atomic_load_32(&tmq->epoch);
1756 1757 1758

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

L
Liu Jicong 已提交
1759
      if (pollRspWrapper->metaRsp.head.epoch == consumerEpoch) {
L
Liu Jicong 已提交
1760
        SMqClientVg* pVg = pollRspWrapper->vgHandle;
S
Shengliang Guan 已提交
1761
        /*printf("vgId:%d, offset %" PRId64 " up to %" PRId64 "\n", pVg->vgId, pVg->currentOffset,
L
Liu Jicong 已提交
1762
         * rspMsg->msg.rspOffset);*/
wmmhello's avatar
wmmhello 已提交
1763
        pVg->currentOffset = pollRspWrapper->metaRsp.rspOffset;
L
Liu Jicong 已提交
1764 1765
        atomic_store_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE);
        // build rsp
L
Liu Jicong 已提交
1766
        SMqMetaRspObj* pRsp = tmqBuildMetaRspFromWrapper(pollRspWrapper);
L
Liu Jicong 已提交
1767 1768 1769
        taosFreeQitem(pollRspWrapper);
        return pRsp;
      } else {
X
Xiaoyu Wang 已提交
1770
        tscDebug("consumer:0x%" PRIx64 " msg discard since epoch mismatch: msg epoch %d, consumer epoch %d",
1771
                 tmq->consumerId, pollRspWrapper->metaRsp.head.epoch, consumerEpoch);
L
Liu Jicong 已提交
1772
        tmqFreeRspWrapper(rspWrapper);
L
Liu Jicong 已提交
1773
        taosFreeQitem(pollRspWrapper);
X
Xiaoyu Wang 已提交
1774
      }
L
Liu Jicong 已提交
1775 1776 1777 1778 1779 1780 1781 1782 1783 1784 1785 1786 1787 1788 1789
    } else if (rspWrapper->tmqRspType == TMQ_MSG_TYPE__TAOSX_RSP) {
      SMqPollRspWrapper* pollRspWrapper = (SMqPollRspWrapper*)rspWrapper;
      /*atomic_sub_fetch_32(&tmq->readyRequest, 1);*/
      int32_t consumerEpoch = atomic_load_32(&tmq->epoch);
      if (pollRspWrapper->taosxRsp.head.epoch == consumerEpoch) {
        SMqClientVg* pVg = pollRspWrapper->vgHandle;
        /*printf("vgId:%d, offset %" PRId64 " up to %" PRId64 "\n", pVg->vgId, pVg->currentOffset,
         * rspMsg->msg.rspOffset);*/
        pVg->currentOffset = pollRspWrapper->taosxRsp.rspOffset;
        atomic_store_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE);
        if (pollRspWrapper->taosxRsp.blockNum == 0) {
          taosFreeQitem(pollRspWrapper);
          rspWrapper = NULL;
          continue;
        }
wmmhello's avatar
wmmhello 已提交
1790

L
Liu Jicong 已提交
1791
        // build rsp
wmmhello's avatar
wmmhello 已提交
1792
        void* pRsp = NULL;
L
Liu Jicong 已提交
1793
        if (pollRspWrapper->taosxRsp.createTableNum == 0) {
wmmhello's avatar
wmmhello 已提交
1794
          pRsp = tmqBuildRspFromWrapper(pollRspWrapper);
L
Liu Jicong 已提交
1795
        } else {
wmmhello's avatar
wmmhello 已提交
1796 1797
          pRsp = tmqBuildTaosxRspFromWrapper(pollRspWrapper);
        }
L
Liu Jicong 已提交
1798 1799 1800
        taosFreeQitem(pollRspWrapper);
        return pRsp;
      } else {
X
Xiaoyu Wang 已提交
1801
        tscDebug("consumer:0x%" PRIx64 " msg discard since epoch mismatch: msg epoch %d, consumer epoch %d",
1802
                 tmq->consumerId, pollRspWrapper->taosxRsp.head.epoch, consumerEpoch);
L
Liu Jicong 已提交
1803
        tmqFreeRspWrapper(rspWrapper);
L
Liu Jicong 已提交
1804 1805
        taosFreeQitem(pollRspWrapper);
      }
X
Xiaoyu Wang 已提交
1806
    } else {
L
fix  
Liu Jicong 已提交
1807
      /*printf("handle ep rsp %d\n", rspMsg->head.mqMsgType);*/
X
Xiaoyu Wang 已提交
1808
      bool reset = false;
L
Liu Jicong 已提交
1809 1810
      tmqHandleNoPollRsp(tmq, rspWrapper, &reset);
      taosFreeQitem(rspWrapper);
X
Xiaoyu Wang 已提交
1811
      if (pollIfReset && reset) {
1812
        tscDebug("consumer:0x%" PRIx64 ", reset and repoll", tmq->consumerId);
1813
        tmqPollImpl(tmq, timeout);
X
Xiaoyu Wang 已提交
1814 1815 1816
      }
    }
  }
H
Haojun Liao 已提交
1817 1818

  tscDebug("consumer:0x%" PRIx64 " handle the rsp completed", tmq->consumerId);
X
Xiaoyu Wang 已提交
1819 1820
}

1821
TAOS_RES* tmq_consumer_poll(tmq_t* tmq, int64_t timeout) {
L
Liu Jicong 已提交
1822 1823
  void*   rspObj;
  int64_t startTime = taosGetTimestampMs();
L
Liu Jicong 已提交
1824

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

1827 1828 1829
#if 0
  tmqHandleAllDelayedTask(tmq);
  tmqPollImpl(tmq, timeout);
1830
  rspObj = tmqHandleAllRsp(tmq, timeout, false);
L
Liu Jicong 已提交
1831 1832
  if (rspObj) {
    return (TAOS_RES*)rspObj;
L
fix  
Liu Jicong 已提交
1833
  }
1834
#endif
X
Xiaoyu Wang 已提交
1835

1836
  // in no topic status, delayed task also need to be processed
L
Liu Jicong 已提交
1837
  if (atomic_load_8(&tmq->status) == TMQ_CONSUMER_STATUS__INIT) {
1838
    tscDebug("consumer:0x%" PRIx64 " poll return since consumer is init", tmq->consumerId);
1839
    taosMsleep(500);  //     sleep for a while
1840 1841 1842
    return NULL;
  }

L
Liu Jicong 已提交
1843 1844 1845
  if (atomic_load_8(&tmq->status) == TMQ_CONSUMER_STATUS__RECOVER) {
    int32_t retryCnt = 0;
    while (TSDB_CODE_MND_CONSUMER_NOT_READY == tmqAskEp(tmq, false)) {
H
Haojun Liao 已提交
1846
      if (retryCnt++ > 40) {
L
Liu Jicong 已提交
1847 1848
        return NULL;
      }
1849

H
Haojun Liao 已提交
1850
      tscDebug("consumer:0x%" PRIx64 " not ready, retry:%d/40 in 500ms", tmq->consumerId, retryCnt);
L
Liu Jicong 已提交
1851 1852 1853 1854
      taosMsleep(500);
    }
  }

X
Xiaoyu Wang 已提交
1855
  while (1) {
L
Liu Jicong 已提交
1856
    tmqHandleAllDelayedTask(tmq);
1857

L
Liu Jicong 已提交
1858
    if (tmqPollImpl(tmq, timeout) < 0) {
1859
      tscDebug("consumer:0x%" PRIx64 " return due to poll error", tmq->consumerId);
L
Liu Jicong 已提交
1860 1861
      /*return NULL;*/
    }
L
Liu Jicong 已提交
1862

1863
    rspObj = tmqHandleAllRsp(tmq, timeout, false);
L
Liu Jicong 已提交
1864
    if (rspObj) {
1865
      tscDebug("consumer:0x%" PRIx64 " return rsp %p", tmq->consumerId, rspObj);
L
Liu Jicong 已提交
1866
      return (TAOS_RES*)rspObj;
L
Liu Jicong 已提交
1867
    } else if (terrno == TSDB_CODE_TQ_NO_COMMITTED_OFFSET) {
1868
      tscDebug("consumer:0x%" PRIx64 " return null since no committed offset", tmq->consumerId);
L
Liu Jicong 已提交
1869
      return NULL;
X
Xiaoyu Wang 已提交
1870
    }
1871

1872
    if (timeout != -1) {
L
Liu Jicong 已提交
1873
      int64_t currentTime = taosGetTimestampMs();
1874 1875 1876
      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 已提交
1877
                 tmq->consumerId, tmq->epoch, startTime, currentTime);
X
Xiaoyu Wang 已提交
1878 1879
        return NULL;
      }
1880
      /*tscInfo("consumer:0x%" PRIx64 ", (epoch %d) wait, start time %" PRId64 ", current time %" PRId64*/
L
Liu Jicong 已提交
1881
      /*", left time %" PRId64,*/
1882 1883
      /*tmq->consumerId, tmq->epoch, startTime, currentTime, (timeout - elapsedTime));*/
      tsem_timewait(&tmq->rspSem, (timeout - elapsedTime));
L
Liu Jicong 已提交
1884 1885
    } else {
      // use tsem_timewait instead of tsem_wait to avoid unexpected stuck
L
Liu Jicong 已提交
1886
      tsem_timewait(&tmq->rspSem, 1000);
X
Xiaoyu Wang 已提交
1887 1888 1889 1890
    }
  }
}

L
Liu Jicong 已提交
1891
int32_t tmq_consumer_close(tmq_t* tmq) {
1892
  if (tmq->status == TMQ_CONSUMER_STATUS__READY) {
L
Liu Jicong 已提交
1893 1894
    int32_t rsp = tmq_commit_sync(tmq, NULL);
    if (rsp != 0) {
L
Liu Jicong 已提交
1895
      return rsp;
1896 1897
    }

L
Liu Jicong 已提交
1898
    int32_t     retryCnt = 0;
1899
    tmq_list_t* lst = tmq_list_new();
L
Liu Jicong 已提交
1900 1901 1902 1903 1904 1905 1906 1907 1908 1909
    while (1) {
      rsp = tmq_subscribe(tmq, lst);
      if (rsp != TSDB_CODE_MND_CONSUMER_NOT_READY || retryCnt > 5) {
        break;
      } else {
        retryCnt++;
        taosMsleep(500);
      }
    }

1910
    tmq_list_destroy(lst);
L
Liu Jicong 已提交
1911
  }
H
Haojun Liao 已提交
1912

1913
  taosRemoveRef(tmqMgmt.rsetId, tmq->refId);
L
Liu Jicong 已提交
1914
  return 0;
1915
}
L
Liu Jicong 已提交
1916

L
Liu Jicong 已提交
1917 1918
const char* tmq_err2str(int32_t err) {
  if (err == 0) {
L
Liu Jicong 已提交
1919
    return "success";
L
Liu Jicong 已提交
1920
  } else if (err == -1) {
L
Liu Jicong 已提交
1921 1922 1923
    return "fail";
  } else {
    return tstrerror(err);
L
Liu Jicong 已提交
1924 1925
  }
}
L
Liu Jicong 已提交
1926

L
Liu Jicong 已提交
1927 1928 1929 1930 1931
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;
1932 1933
  } else if (TD_RES_TMQ_METADATA(res)) {
    return TMQ_RES_METADATA;
L
Liu Jicong 已提交
1934 1935 1936 1937 1938
  } else {
    return TMQ_RES_INVALID;
  }
}

L
Liu Jicong 已提交
1939
const char* tmq_get_topic_name(TAOS_RES* res) {
L
Liu Jicong 已提交
1940 1941
  if (TD_RES_TMQ(res)) {
    SMqRspObj* pRspObj = (SMqRspObj*)res;
L
Liu Jicong 已提交
1942
    return strchr(pRspObj->topic, '.') + 1;
L
Liu Jicong 已提交
1943 1944 1945
  } else if (TD_RES_TMQ_META(res)) {
    SMqMetaRspObj* pMetaRspObj = (SMqMetaRspObj*)res;
    return strchr(pMetaRspObj->topic, '.') + 1;
1946 1947 1948
  } else if (TD_RES_TMQ_METADATA(res)) {
    SMqTaosxRspObj* pRspObj = (SMqTaosxRspObj*)res;
    return strchr(pRspObj->topic, '.') + 1;
L
Liu Jicong 已提交
1949 1950 1951 1952 1953
  } else {
    return NULL;
  }
}

L
Liu Jicong 已提交
1954 1955 1956 1957
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 已提交
1958 1959 1960
  } else if (TD_RES_TMQ_META(res)) {
    SMqMetaRspObj* pMetaRspObj = (SMqMetaRspObj*)res;
    return strchr(pMetaRspObj->db, '.') + 1;
1961 1962 1963
  } else if (TD_RES_TMQ_METADATA(res)) {
    SMqTaosxRspObj* pRspObj = (SMqTaosxRspObj*)res;
    return strchr(pRspObj->db, '.') + 1;
L
Liu Jicong 已提交
1964 1965 1966 1967 1968
  } else {
    return NULL;
  }
}

L
Liu Jicong 已提交
1969 1970 1971 1972
int32_t tmq_get_vgroup_id(TAOS_RES* res) {
  if (TD_RES_TMQ(res)) {
    SMqRspObj* pRspObj = (SMqRspObj*)res;
    return pRspObj->vgId;
L
Liu Jicong 已提交
1973 1974 1975
  } else if (TD_RES_TMQ_META(res)) {
    SMqMetaRspObj* pMetaRspObj = (SMqMetaRspObj*)res;
    return pMetaRspObj->vgId;
1976
  } else if (TD_RES_TMQ_METADATA(res)) {
1977 1978
    SMqTaosxRspObj* pRspObj = (SMqTaosxRspObj*)res;
    return pRspObj->vgId;
L
Liu Jicong 已提交
1979 1980 1981 1982
  } else {
    return -1;
  }
}
L
Liu Jicong 已提交
1983 1984 1985 1986 1987 1988 1989 1990

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;
    }
1991
    return (const char*)taosArrayGetP(pRspObj->rsp.blockTbName, pRspObj->resIter);
1992 1993
  } else if (TD_RES_TMQ_METADATA(res)) {
    SMqTaosxRspObj* pRspObj = (SMqTaosxRspObj*)res;
L
Liu Jicong 已提交
1994 1995 1996
    if (!pRspObj->rsp.withTbName || pRspObj->rsp.blockTbName == NULL || pRspObj->resIter < 0 ||
        pRspObj->resIter >= pRspObj->rsp.blockNum) {
      return NULL;
1997
    }
L
Liu Jicong 已提交
1998 1999
    return (const char*)taosArrayGetP(pRspObj->rsp.blockTbName, pRspObj->resIter);
  }
L
Liu Jicong 已提交
2000 2001
  return NULL;
}
2002

L
Liu Jicong 已提交
2003
void tmq_commit_async(tmq_t* tmq, const TAOS_RES* msg, tmq_commit_cb* cb, void* param) {
L
Liu Jicong 已提交
2004
  tmqCommitInner(tmq, msg, 0, 1, cb, param);
L
Liu Jicong 已提交
2005 2006
}

2007
int32_t tmq_commit_sync(tmq_t* tmq, const TAOS_RES* msg) {
L
Liu Jicong 已提交
2008
  return tmqCommitInner(tmq, msg, 0, 0, NULL, NULL);
2009
}
2010 2011 2012 2013 2014 2015 2016 2017 2018 2019 2020 2021 2022 2023 2024 2025 2026 2027 2028 2029 2030 2031 2032 2033 2034 2035 2036 2037 2038 2039 2040 2041 2042 2043 2044 2045 2046 2047 2048 2049 2050 2051 2052 2053 2054 2055 2056 2057 2058 2059 2060 2061 2062 2063 2064 2065 2066 2067 2068 2069 2070 2071 2072 2073 2074 2075 2076 2077 2078 2079 2080 2081 2082 2083 2084 2085 2086 2087 2088 2089 2090 2091 2092 2093 2094 2095 2096 2097 2098 2099 2100 2101 2102 2103 2104 2105 2106 2107 2108 2109 2110 2111 2112 2113 2114

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) {
  tmq_t* tmq = taosAcquireRef(tmqMgmt.rsetId, pParamSet->refId);
  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);
2115
    } else if (!pParamSet->automatic && pParamSet->userCb) { // sem post
2116 2117
      pParamSet->userCb(tmq, pParamSet->rspErr, pParamSet->userParam);
    }
2118

2119 2120 2121 2122 2123 2124 2125 2126 2127 2128
    taosMemoryFree(pParamSet);
  } else {
    tsem_post(&pParamSet->rspSem);
  }

#if 0
  taosArrayDestroyP(pParamSet->successfulOffsets, taosMemoryFree);
    taosArrayDestroyP(pParamSet->failedOffsets, taosMemoryFree);
#endif
  return 0;
2129 2130 2131 2132 2133 2134 2135 2136 2137 2138 2139
}

void tmqCommitRspCountDown(SMqCommitCbParamSet* pParamSet, int64_t consumerId, const char* pTopic, int32_t vgId) {
  int32_t waitingRspNum = atomic_sub_fetch_32(&pParamSet->waitingRspNum, 1);
  tscDebug("consumer:0x%" PRIx64 " topic:%s vgId:%d commit-rsp received, remain:%d", consumerId, pTopic, vgId,
           waitingRspNum);

  if (waitingRspNum == 0) {
    tmqCommitDone(pParamSet);
  }
}