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

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

H
Haojun Liao 已提交
27 28
#define VG_POLL_IGNORE_TICK 100

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

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

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

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

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

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

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

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

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

L
Liu Jicong 已提交
97
  // timer
98 99
  tmr_h hbLiveTimer;
  tmr_h epTimer;
L
Liu Jicong 已提交
100 101 102
  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;
L
Liu Jicong 已提交
132 133
  STqOffsetVal committedOffset;
  STqOffsetVal currentOffset;
H
Haojun Liao 已提交
134 135 136 137 138
  int32_t      vgId;
  int32_t      vgStatus;
  int32_t      vgSkipCnt;
  int32_t      vgIgnoreCnt; // once empty block is received, idle for ignoreCnt then start to poll data
  SEpSet       epSet;
139 140
} 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;
H
Haojun Liao 已提交
153
  uint64_t        reqId;
L
Liu Jicong 已提交
154
  union {
L
Liu Jicong 已提交
155 156
    SMqDataRsp dataRsp;
    SMqMetaRsp metaRsp;
L
Liu Jicong 已提交
157
    STaosxRsp  taosxRsp;
L
Liu Jicong 已提交
158
  };
L
Liu Jicong 已提交
159 160
} SMqPollRspWrapper;

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

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

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

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

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

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

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

223
  conf->withTbName = false;
L
Liu Jicong 已提交
224
  conf->autoCommit = true;
L
Liu Jicong 已提交
225
  conf->autoCommitInterval = 5000;
226
  conf->resetOffset = TMQ_OFFSET__RESET_EARLIEAST;
227
  conf->hbBgEnable = true;
228

229 230 231
  return conf;
}

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

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

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

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

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

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

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

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

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

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

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

L
Liu Jicong 已提交
351
  return TMQ_CONF_UNKNOWN;
352 353 354
}

tmq_list_t* tmq_list_new() {
L
Liu Jicong 已提交
355
  return (tmq_list_t*)taosArrayInit(0, sizeof(void*));
356 357
}

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

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

L
Liu Jicong 已提交
374 375 376 377 378 379 380 381 382 383
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;
}

384 385 386 387 388
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;

389
  for(int32_t i = 0; i < numOfTopics; ++i) {
390 391
    SMqClientTopic* pTopic = taosArrayGet(pTopicList, i);
    if (strcmp(pTopic->topicName, pName) != 0) {
392 393 394
      continue;
    }

395 396
    *numOfVgroups = taosArrayGetSize(pTopic->vgs);
    for (int32_t j = 0; j < (*numOfVgroups); ++j) {
397
      SMqClientVg* pClientVg = taosArrayGet(pTopic->vgs, j);
398 399 400
      if (pClientVg->vgId == vgId) {
        *index = j;
        return pClientVg;
401 402
      }
    }
L
Liu Jicong 已提交
403
  }
404 405

  return NULL;
L
Liu Jicong 已提交
406
}
407

H
Haojun Liao 已提交
408
static int32_t tmqCommitCb(void* param, SDataBuf* pBuf, int32_t code) {
409
  SMqCommitCbParam*    pParam = (SMqCommitCbParam*)param;
410
  SMqCommitCbParamSet* pParamSet = (SMqCommitCbParamSet*)pParam->params;
H
Haojun Liao 已提交
411

412
  if (code != TSDB_CODE_SUCCESS) { // if commit offset failed, let's try again
H
Haojun Liao 已提交
413
    taosThreadMutexLock(&pParam->pTmq->lock);
414 415 416 417 418 419 420 421 422 423 424 425 426 427 428 429 430
    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);
      }
    }

H
Haojun Liao 已提交
431
    taosThreadMutexUnlock(&pParam->pTmq->lock);
432 433 434 435 436 437 438 439 440 441 442 443 444 445 446 447 448 449 450 451 452 453 454 455 456 457 458 459 460 461

    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
  taosThreadMutexLock(&pParam->pTmq->lock);
  tmq_t* pTmq = pParam->pTmq;
  int32_t index = 0, numOfVgroups = 0;

  SMqClientVg* pVg = foundClientVg(pTmq->clientTopics, pParam->topicName, pParam->vgId, &index, &numOfVgroups);
  if (pVg == NULL) {
    tscDebug("consumer:0x%" PRIx64 " subKey:%s vgId:%d has been transferred to other consumer, ordinal:%d/%d",
             pParam->pTmq->consumerId, pParam->pOffset->subKey, pParam->vgId, index + 1, numOfVgroups);
  } else { // update the epset if needed
    if (pBuf->pEpSet != NULL) {
      SEp* pEp = GET_ACTIVE_EP(pBuf->pEpSet);
      SEp* pOld = GET_ACTIVE_EP(&(pVg->epSet));

      tscDebug("consumer:0x%" PRIx64 " subKey:%s update the epset vgId:%d, ep:%s:%d, old ep:%s:%d, ordinal:%d/%d",
               pTmq->consumerId, pParam->pOffset->subKey, pParam->vgId, pEp->fqdn, pEp->port, pOld->fqdn, pOld->port,
               index + 1, numOfVgroups);

      pVg->epSet = *pBuf->pEpSet;
    }

H
Haojun Liao 已提交
462 463
    tscDebug("consumer:0x%" PRIx64 " subKey:%s vgId:%d, commit offset success. ordinal:%d/%d", pTmq->consumerId,
             pParam->pOffset->subKey, pParam->vgId, index + 1, numOfVgroups);
464 465
  }

466 467
  taosThreadMutexUnlock(&pParam->pTmq->lock);

L
Liu Jicong 已提交
468
  taosMemoryFree(pParam->pOffset);
L
Liu Jicong 已提交
469
  taosMemoryFree(pBuf->pData);
dengyihao's avatar
dengyihao 已提交
470
  taosMemoryFree(pBuf->pEpSet);
L
Liu Jicong 已提交
471

472
  tmqCommitRspCountDown(pParamSet, pParam->pTmq->consumerId, pParam->topicName, pParam->vgId);
473 474 475
  return 0;
}

476 477
static int32_t doSendCommitMsg(tmq_t* tmq, SMqClientVg* pVg, const char* pTopicName, SMqCommitCbParamSet* pParamSet,
                               int32_t index, int32_t totalVgroups) {
L
Liu Jicong 已提交
478 479 480 481 482
  STqOffset* pOffset = taosMemoryCalloc(1, sizeof(STqOffset));
  if (pOffset == NULL) {
    terrno = TSDB_CODE_OUT_OF_MEMORY;
    return -1;
  }
483

L
Liu Jicong 已提交
484
  pOffset->val = pVg->currentOffset;
485

L
Liu Jicong 已提交
486 487 488
  int32_t groupLen = strlen(tmq->groupId);
  memcpy(pOffset->subKey, tmq->groupId, groupLen);
  pOffset->subKey[groupLen] = TMQ_SEPARATOR;
H
Haojun Liao 已提交
489
  strcpy(pOffset->subKey + groupLen + 1, pTopicName);
L
Liu Jicong 已提交
490

L
Liu Jicong 已提交
491 492 493 494 495 496
  int32_t len;
  int32_t code;
  tEncodeSize(tEncodeSTqOffset, pOffset, len, code);
  if (code < 0) {
    return -1;
  }
497

L
Liu Jicong 已提交
498
  void* buf = taosMemoryCalloc(1, sizeof(SMsgHead) + len);
L
Liu Jicong 已提交
499 500 501 502
  if (buf == NULL) {
    taosMemoryFree(pOffset);
    return -1;
  }
503

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

L
Liu Jicong 已提交
506 507 508 509 510
  void* abuf = POINTER_SHIFT(buf, sizeof(SMsgHead));

  SEncoder encoder;
  tEncoderInit(&encoder, abuf, len);
  tEncodeSTqOffset(&encoder, pOffset);
L
Liu Jicong 已提交
511
  tEncoderClear(&encoder);
L
Liu Jicong 已提交
512 513

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

L
Liu Jicong 已提交
521 522
  pParam->params = pParamSet;
  pParam->pOffset = pOffset;
H
Haojun Liao 已提交
523 524 525
  pParam->vgId = pVg->vgId;
  pParam->pTmq = tmq;

H
Haojun Liao 已提交
526
  tstrncpy(pParam->topicName, pTopicName, tListLen(pParam->topicName));
L
Liu Jicong 已提交
527 528 529 530

  // build send info
  SMsgSendInfo* pMsgSendInfo = taosMemoryCalloc(1, sizeof(SMsgSendInfo));
  if (pMsgSendInfo == NULL) {
L
Liu Jicong 已提交
531
    taosMemoryFree(pOffset);
L
Liu Jicong 已提交
532 533
    taosMemoryFree(buf);
    taosMemoryFree(pParam);
L
Liu Jicong 已提交
534 535
    return -1;
  }
536

L
Liu Jicong 已提交
537 538 539 540 541 542
  pMsgSendInfo->msgInfo = (SDataBuf){
      .pData = buf,
      .len = sizeof(SMsgHead) + len,
      .handle = NULL,
  };

H
Haojun Liao 已提交
543
  SEp* pEp = GET_ACTIVE_EP(&pVg->epSet);
544 545 546
  tscDebug("consumer:0x%" PRIx64 " topic:%s on vgId:%d offset:%" PRId64 " prev:%" PRId64 ", ep:%s:%d, ordinal:%d/%d",
           tmq->consumerId, pOffset->subKey, pVg->vgId, pOffset->val.version, pVg->committedOffset.version, pEp->fqdn,
           pEp->port, index + 1, totalVgroups);
L
Liu Jicong 已提交
547 548 549 550

  pMsgSendInfo->requestId = generateRequestId();
  pMsgSendInfo->requestObjRefId = 0;
  pMsgSendInfo->param = pParam;
L
Liu Jicong 已提交
551
  pMsgSendInfo->paramFreeFp = taosMemoryFree;
552
  pMsgSendInfo->fp = tmqCommitCb;
L
Liu Jicong 已提交
553
  pMsgSendInfo->msgType = TDMT_VND_TMQ_COMMIT_OFFSET;
L
Liu Jicong 已提交
554

L
Liu Jicong 已提交
555 556 557
  atomic_add_fetch_32(&pParamSet->waitingRspNum, 1);
  atomic_add_fetch_32(&pParamSet->totalRspNum, 1);

L
Liu Jicong 已提交
558 559 560 561 562
  int64_t transporterId = 0;
  asyncSendMsgToServer(tmq->pTscObj->pAppInfo->pTransporter, &pVg->epSet, &transporterId, pMsgSendInfo);
  return 0;
}

H
Haojun Liao 已提交
563
static int32_t tmqCommitMsgImpl(tmq_t* tmq, const TAOS_RES* msg, int8_t async, tmq_commit_cb* userCb, void* userParam) {
L
Liu Jicong 已提交
564 565 566 567 568 569 570 571 572 573
  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 已提交
574
  } else if (TD_RES_TMQ_METADATA(msg)) {
575 576 577
    SMqTaosxRspObj* pRspObj = (SMqTaosxRspObj*)msg;
    topic = pRspObj->topic;
    vgId = pRspObj->vgId;
L
Liu Jicong 已提交
578 579 580 581 582 583 584 585 586
  } 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 已提交
587

588 589
  pParamSet->refId = tmq->refId;
  pParamSet->epoch = tmq->epoch;
L
Liu Jicong 已提交
590 591 592 593 594 595
  pParamSet->automatic = 0;
  pParamSet->async = async;
  pParamSet->userCb = userCb;
  pParamSet->userParam = userParam;
  tsem_init(&pParamSet->rspSem, 0, 0);

L
Liu Jicong 已提交
596 597
  int32_t code = -1;

H
Haojun Liao 已提交
598
  taosThreadMutexLock(&tmq->lock);
L
Liu Jicong 已提交
599 600
  for (int32_t i = 0; i < taosArrayGetSize(tmq->clientTopics); i++) {
    SMqClientTopic* pTopic = taosArrayGet(tmq->clientTopics, i);
601 602 603 604 605 606
    if (strcmp(pTopic->topicName, topic) != 0) {
      continue;
    }

    int32_t numOfVgroups = taosArrayGetSize(pTopic->vgs);
    for (int32_t j = 0; j < numOfVgroups; j++) {
L
Liu Jicong 已提交
607
      SMqClientVg* pVg = taosArrayGet(pTopic->vgs, j);
H
Haojun Liao 已提交
608 609 610
      if (pVg->vgId != vgId) {
        continue;
      }
L
Liu Jicong 已提交
611

L
Liu Jicong 已提交
612
      if (pVg->currentOffset.type > 0 && !tOffsetEqual(&pVg->currentOffset, &pVg->committedOffset)) {
613
        if (doSendCommitMsg(tmq, pVg, pTopic->topicName, pParamSet, j, numOfVgroups) < 0) {
L
Liu Jicong 已提交
614 615
          tsem_destroy(&pParamSet->rspSem);
          taosMemoryFree(pParamSet);
L
Liu Jicong 已提交
616
          goto FAIL;
L
Liu Jicong 已提交
617
        }
L
Liu Jicong 已提交
618
        goto HANDLE_RSP;
L
Liu Jicong 已提交
619 620
      }
    }
L
Liu Jicong 已提交
621
  }
L
Liu Jicong 已提交
622

L
Liu Jicong 已提交
623 624 625 626
HANDLE_RSP:
  if (pParamSet->totalRspNum == 0) {
    tsem_destroy(&pParamSet->rspSem);
    taosMemoryFree(pParamSet);
H
Haojun Liao 已提交
627
    taosThreadMutexUnlock(&tmq->lock);
L
Liu Jicong 已提交
628 629 630
    return 0;
  }

L
Liu Jicong 已提交
631 632 633 634
  if (!async) {
    tsem_wait(&pParamSet->rspSem);
    code = pParamSet->rspErr;
    tsem_destroy(&pParamSet->rspSem);
L
Liu Jicong 已提交
635
    taosMemoryFree(pParamSet);
H
Haojun Liao 已提交
636
    taosThreadMutexUnlock(&tmq->lock);
L
Liu Jicong 已提交
637 638 639 640 641 642
    return code;
  } else {
    code = 0;
  }

FAIL:
H
Haojun Liao 已提交
643
  taosThreadMutexUnlock(&tmq->lock);
L
Liu Jicong 已提交
644 645 646
  if (code != 0 && async) {
    userCb(tmq, code, userParam);
  }
H
Haojun Liao 已提交
647

L
Liu Jicong 已提交
648 649 650
  return 0;
}

L
Liu Jicong 已提交
651 652
static int32_t tmqCommitConsumerImpl(tmq_t* tmq, int8_t automatic, int8_t async, tmq_commit_cb* userCb,
                                     void* userParam) {
L
Liu Jicong 已提交
653 654
  int32_t code = -1;

655 656
  SMqCommitCbParamSet* pParamSet = taosMemoryCalloc(1, sizeof(SMqCommitCbParamSet));
  if (pParamSet == NULL) {
L
Liu Jicong 已提交
657 658 659 660 661 662 663 664
    code = TSDB_CODE_OUT_OF_MEMORY;
    if (async) {
      if (automatic) {
        tmq->commitCb(tmq, code, tmq->commitCbUserParam);
      } else {
        userCb(tmq, code, userParam);
      }
    }
665 666
    return -1;
  }
667 668 669 670

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

671 672 673 674 675 676
  pParamSet->automatic = automatic;
  pParamSet->async = async;
  pParamSet->userCb = userCb;
  pParamSet->userParam = userParam;
  tsem_init(&pParamSet->rspSem, 0, 0);

677 678 679
  // init as 1 to prevent concurrency issue
  pParamSet->waitingRspNum = 1;

H
Haojun Liao 已提交
680
  taosThreadMutexLock(&tmq->lock);
681
  int32_t numOfTopics = taosArrayGetSize(tmq->clientTopics);
X
Xiaoyu Wang 已提交
682
  tscDebug("consumer:0x%" PRIx64 " start to commit offset for %d topics", tmq->consumerId, numOfTopics);
683 684

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

688 689
    tscDebug("consumer:0x%" PRIx64 " commit offset for topics:%s, numOfVgs:%d", tmq->consumerId, pTopic->topicName,
             numOfVgroups);
690
    for (int32_t j = 0; j < numOfVgroups; j++) {
691 692 693 694 695 696 697 698
      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 已提交
699 700
          continue;
        }
H
Haojun Liao 已提交
701 702 703

        // update the offset value.
        pVg->committedOffset = pVg->currentOffset;
704
      } else {
705 706
        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);
707 708 709 710
      }
    }
  }

H
Haojun Liao 已提交
711
  tscDebug("consumer:0x%" PRIx64 " total commit:%d for %d topics", tmq->consumerId, pParamSet->waitingRspNum - 1,
712
           numOfTopics);
H
Haojun Liao 已提交
713 714
  taosThreadMutexUnlock(&tmq->lock);

L
Liu Jicong 已提交
715
  // no request is sent
L
Liu Jicong 已提交
716 717 718 719 720 721
  if (pParamSet->totalRspNum == 0) {
    tsem_destroy(&pParamSet->rspSem);
    taosMemoryFree(pParamSet);
    return 0;
  }

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

725 726 727 728
  if (!async) {
    tsem_wait(&pParamSet->rspSem);
    code = pParamSet->rspErr;
    tsem_destroy(&pParamSet->rspSem);
729
    taosMemoryFree(pParamSet);
L
Liu Jicong 已提交
730
#if 0
731 732
    taosArrayDestroyP(pParamSet->successfulOffsets, taosMemoryFree);
    taosArrayDestroyP(pParamSet->failedOffsets, taosMemoryFree);
L
Liu Jicong 已提交
733
#endif
L
Liu Jicong 已提交
734
  }
735

L
Liu Jicong 已提交
736 737 738
  return code;
}

739 740
static int32_t tmqCommitInner(tmq_t* tmq, const TAOS_RES* msg, int8_t automatic, int8_t async, tmq_commit_cb* userCb,
                              void* userParam) {
741
  if (msg) { // user invoked commit?
L
Liu Jicong 已提交
742
    return tmqCommitMsgImpl(tmq, msg, async, userCb, userParam);
743
  } else {  // this for auto commit
L
Liu Jicong 已提交
744 745
    return tmqCommitConsumerImpl(tmq, automatic, async, userCb, userParam);
  }
746 747
}

748
void tmqAssignAskEpTask(void* param, void* tmrId) {
749 750 751
  int64_t refId = *(int64_t*)param;
  tmq_t*  tmq = taosAcquireRef(tmqMgmt.rsetId, refId);
  if (tmq != NULL) {
S
Shengliang Guan 已提交
752
    int8_t* pTaskType = taosAllocateQitem(sizeof(int8_t), DEF_QITEM, 0);
753 754 755 756 757
    *pTaskType = TMQ_DELAYED_TASK__ASK_EP;
    taosWriteQitem(tmq->delayedTask, pTaskType);
    tsem_post(&tmq->rspSem);
  }
  taosMemoryFree(param);
L
Liu Jicong 已提交
758 759 760
}

void tmqAssignDelayedCommitTask(void* param, void* tmrId) {
761 762 763
  int64_t refId = *(int64_t*)param;
  tmq_t*  tmq = taosAcquireRef(tmqMgmt.rsetId, refId);
  if (tmq != NULL) {
S
Shengliang Guan 已提交
764
    int8_t* pTaskType = taosAllocateQitem(sizeof(int8_t), DEF_QITEM, 0);
765 766 767 768 769
    *pTaskType = TMQ_DELAYED_TASK__COMMIT;
    taosWriteQitem(tmq->delayedTask, pTaskType);
    tsem_post(&tmq->rspSem);
  }
  taosMemoryFree(param);
L
Liu Jicong 已提交
770 771 772
}

void tmqAssignDelayedReportTask(void* param, void* tmrId) {
773 774 775
  int64_t refId = *(int64_t*)param;
  tmq_t*  tmq = taosAcquireRef(tmqMgmt.rsetId, refId);
  if (tmq != NULL) {
S
Shengliang Guan 已提交
776
    int8_t* pTaskType = taosAllocateQitem(sizeof(int8_t), DEF_QITEM, 0);
777 778 779 780 781
    *pTaskType = TMQ_DELAYED_TASK__REPORT;
    taosWriteQitem(tmq->delayedTask, pTaskType);
    tsem_post(&tmq->rspSem);
  }
  taosMemoryFree(param);
L
Liu Jicong 已提交
782 783
}

784
int32_t tmqHbCb(void* param, SDataBuf* pMsg, int32_t code) {
dengyihao's avatar
dengyihao 已提交
785 786 787 788
  if (pMsg) {
    taosMemoryFree(pMsg->pData);
    taosMemoryFree(pMsg->pEpSet);
  }
789 790 791 792
  return 0;
}

void tmqSendHbReq(void* param, void* tmrId) {
793 794 795
  int64_t refId = *(int64_t*)param;
  tmq_t*  tmq = taosAcquireRef(tmqMgmt.rsetId, refId);
  if (tmq == NULL) {
L
Liu Jicong 已提交
796
    taosMemoryFree(param);
797 798
    return;
  }
D
dapan1121 已提交
799 800 801 802 803

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

L
Liu Jicong 已提交
804
  int32_t tlen = tSerializeSMqHbReq(NULL, 0, &req);
D
dapan1121 已提交
805 806 807 808
  if (tlen < 0) {
    tscError("tSerializeSMqHbReq failed");
    return;
  }
L
Liu Jicong 已提交
809
  void* pReq = taosMemoryCalloc(1, tlen);
D
dapan1121 已提交
810 811 812 813 814 815 816 817 818
  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;
  }
819 820 821 822

  SMsgSendInfo* sendInfo = taosMemoryCalloc(1, sizeof(SMsgSendInfo));
  if (sendInfo == NULL) {
    taosMemoryFree(pReq);
L
Liu Jicong 已提交
823
    goto OVER;
824 825 826
  }
  sendInfo->msgInfo = (SDataBuf){
      .pData = pReq,
D
dapan1121 已提交
827
      .len = tlen,
828 829 830 831 832 833 834
      .handle = NULL,
  };

  sendInfo->requestId = generateRequestId();
  sendInfo->requestObjRefId = 0;
  sendInfo->param = NULL;
  sendInfo->fp = tmqHbCb;
L
Liu Jicong 已提交
835
  sendInfo->msgType = TDMT_MND_TMQ_HB;
836 837 838 839 840 841 842

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

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

OVER:
843
  taosTmrReset(tmqSendHbReq, 1000, param, tmqMgmt.timer, &tmq->hbLiveTimer);
844 845
}

846
int32_t tmqHandleAllDelayedTask(tmq_t* pTmq) {
L
Liu Jicong 已提交
847
  STaosQall* qall = taosAllocateQall();
848
  taosReadAllQitems(pTmq->delayedTask, qall);
L
Liu Jicong 已提交
849

850 851 852 853
  if (qall->numOfItems == 0) {
    taosFreeQall(qall);
    return TSDB_CODE_SUCCESS;
  }
854

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

859
  while (pTaskType != NULL) {
860
    if (*pTaskType == TMQ_DELAYED_TASK__ASK_EP) {
861
      tmqAskEp(pTmq, true);
862 863

      int64_t* pRefId = taosMemoryMalloc(sizeof(int64_t));
864
      *pRefId = pTmq->refId;
865

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

      int64_t* pRefId = taosMemoryMalloc(sizeof(int64_t));
872
      *pRefId = pTmq->refId;
873

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

L
Liu Jicong 已提交
880
    taosFreeQitem(pTaskType);
881
    taosGetQitem(qall, (void**)&pTaskType);
L
Liu Jicong 已提交
882
  }
883

L
Liu Jicong 已提交
884 885 886 887
  taosFreeQall(qall);
  return 0;
}

L
Liu Jicong 已提交
888 889 890 891 892 893 894 895 896 897 898 899 900 901 902 903 904 905 906 907 908 909 910 911 912 913 914
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 已提交
915
void tmqClearUnhandleMsg(tmq_t* tmq) {
L
Liu Jicong 已提交
916
  SMqRspWrapper* rspWrapper = NULL;
L
Liu Jicong 已提交
917
  while (1) {
L
Liu Jicong 已提交
918 919 920 921 922
    taosGetQitem(tmq->qall, (void**)&rspWrapper);
    if (rspWrapper) {
      tmqFreeRspWrapper(rspWrapper);
      taosFreeQitem(rspWrapper);
    } else {
L
Liu Jicong 已提交
923
      break;
L
Liu Jicong 已提交
924
    }
L
Liu Jicong 已提交
925 926
  }

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

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

  taosMemoryFree(pMsg->pEpSet);
L
Liu Jicong 已提交
945 946 947
  tsem_post(&pParam->rspSem);
  return 0;
}
948

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

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

978 979
void tmqFreeImpl(void* handle) {
  tmq_t* tmq = (tmq_t*)handle;
L
Liu Jicong 已提交
980

981
  // TODO stop timer
L
Liu Jicong 已提交
982 983 984 985
  if (tmq->mqueue) {
    tmqClearUnhandleMsg(tmq);
    taosCloseQueue(tmq->mqueue);
  }
L
Liu Jicong 已提交
986

H
Haojun Liao 已提交
987 988 989 990 991
  if (tmq->delayedTask) {
    taosCloseQueue(tmq->delayedTask);
  }

  taosFreeQall(tmq->qall);
992
  tsem_destroy(&tmq->rspSem);
H
Haojun Liao 已提交
993
  taosThreadMutexDestroy(&tmq->lock);
L
Liu Jicong 已提交
994

995 996 997
  int32_t sz = taosArrayGetSize(tmq->clientTopics);
  for (int32_t i = 0; i < sz; i++) {
    SMqClientTopic* pTopic = taosArrayGet(tmq->clientTopics, i);
L
Liu Jicong 已提交
998
    taosMemoryFreeClear(pTopic->schema.pSchema);
999 1000
    taosArrayDestroy(pTopic->vgs);
  }
H
Haojun Liao 已提交
1001

1002 1003 1004
  taosArrayDestroy(tmq->clientTopics);
  taos_close_internal(tmq->pTscObj);
  taosMemoryFree(tmq);
L
Liu Jicong 已提交
1005 1006
}

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

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

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

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

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

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

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

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

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

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

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

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

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

1101 1102 1103 1104 1105 1106
  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 已提交
1107

1108
  return pTmq;
1109

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

D
dapan1121 已提交
1244
int32_t tmqPollCb(void* param, SDataBuf* pMsg, int32_t code) {
X
Xiaoyu Wang 已提交
1245 1246
  SMqPollCbParam* pParam = (SMqPollCbParam*)param;
  SMqClientVg*    pVg = pParam->pVg;
L
Liu Jicong 已提交
1247
  SMqClientTopic* pTopic = pParam->pTopic;
1248 1249 1250 1251 1252 1253

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

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

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

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

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

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

L
Liu Jicong 已提交
1285 1286 1287 1288
      pRspWrapper->tmqRspType = TMQ_MSG_TYPE__END_RSP;
      taosWriteQitem(tmq->mqueue, pRspWrapper);
      tsem_post(&tmq->rspSem);
    }
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);
L
Liu Jicong 已提交
1301
    taosMemoryFree(pMsg->pData);
dengyihao's avatar
dengyihao 已提交
1302
    taosMemoryFree(pMsg->pEpSet);
X
Xiaoyu Wang 已提交
1303 1304 1305 1306
    return 0;
  }

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

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

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

L
Liu Jicong 已提交
1322
  pRspWrapper->tmqRspType = rspType;
L
Liu Jicong 已提交
1323 1324
  pRspWrapper->vgHandle = pVg;
  pRspWrapper->topicHandle = pTopic;
H
Haojun Liao 已提交
1325
  pRspWrapper->reqId = requestId;
L
Liu Jicong 已提交
1326

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

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

L
Liu Jicong 已提交
1338
  } else if (rspType == TMQ_MSG_TYPE__POLL_META_RSP) {
1339 1340 1341 1342
    SDecoder decoder;
    tDecoderInit(&decoder, POINTER_SHIFT(pMsg->pData, sizeof(SMqRspHead)), pMsg->len - sizeof(SMqRspHead));
    tDecodeSMqMetaRsp(&decoder, &pRspWrapper->metaRsp);
    tDecoderClear(&decoder);
L
Liu Jicong 已提交
1343
    memcpy(&pRspWrapper->metaRsp, pMsg->pData, sizeof(SMqRspHead));
L
Liu Jicong 已提交
1344 1345 1346 1347 1348 1349
  } 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 已提交
1350 1351
  } else { // invalid rspType
    tscError("consumer:0x%"PRIx64" invalid rsp msg received, type:%d ignored", tmq->consumerId, rspType);
L
Liu Jicong 已提交
1352
  }
L
Liu Jicong 已提交
1353

L
Liu Jicong 已提交
1354
  taosMemoryFree(pMsg->pData);
dengyihao's avatar
dengyihao 已提交
1355
  taosMemoryFree(pMsg->pEpSet);
H
Haojun Liao 已提交
1356
  taosWriteQitem(tmq->mqueue, pRspWrapper);
L
Liu Jicong 已提交
1357

H
Haojun Liao 已提交
1358 1359 1360
  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);

1361
  tsem_post(&tmq->rspSem);
L
Liu Jicong 已提交
1362
  return 0;
H
Haojun Liao 已提交
1363

L
fix txn  
Liu Jicong 已提交
1364
CREATE_MSG_FAIL:
L
Liu Jicong 已提交
1365
  if (epoch == tmq->epoch) {
L
Liu Jicong 已提交
1366 1367
    atomic_store_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE);
  }
H
Haojun Liao 已提交
1368

1369
  tsem_post(&tmq->rspSem);
L
Liu Jicong 已提交
1370
  return -1;
1371 1372
}

H
Haojun Liao 已提交
1373 1374 1375 1376 1377 1378 1379 1380 1381 1382 1383 1384 1385 1386 1387 1388 1389
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 已提交
1390 1391

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

H
Haojun Liao 已提交
1394 1395 1396 1397 1398 1399 1400 1401 1402 1403 1404 1405
    STqOffsetVal  offsetNew = {.type = tmq->resetOffsetCfg};
    if (pOffset != NULL) {
      offsetNew = *pOffset;
    }

    SMqClientVg clientVg = {
        .pollCnt = 0,
        .currentOffset = offsetNew,
        .vgId = pVgEp->vgId,
        .epSet = pVgEp->epSet,
        .vgStatus = TMQ_VG_STATUS__IDLE,
        .vgSkipCnt = 0,
H
Haojun Liao 已提交
1406
        .vgIgnoreCnt = 0,
H
Haojun Liao 已提交
1407 1408 1409 1410 1411 1412 1413 1414 1415 1416 1417 1418 1419 1420 1421 1422
    };

    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) {
1423 1424
  bool set = false;

1425
  int32_t topicNumCur = taosArrayGetSize(tmq->clientTopics);
1426
  int32_t topicNumGet = taosArrayGetSize(pRsp->topics);
1427

X
Xiaoyu Wang 已提交
1428 1429
  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",
1430
           tmq->consumerId, tmq->epoch, epoch, topicNumGet, topicNumCur);
1431 1432 1433 1434 1435 1436

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

H
Haojun Liao 已提交
1437 1438
  SHashObj* pVgOffsetHashMap = taosHashInit(64, MurmurHash3_32, false, HASH_NO_LOCK);
  if (pVgOffsetHashMap == NULL) {
1439 1440 1441
    taosArrayDestroy(newTopics);
    return false;
  }
1442

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

L
Liu Jicong 已提交
1454
        char buf[80];
L
Liu Jicong 已提交
1455
        tFormatOffset(buf, 80, &pVgCur->currentOffset);
H
Haojun Liao 已提交
1456
        tscDebug("consumer:0x%" PRIx64 ", epoch:%d vgId:%d vgKey:%s, offset:%s", tmq->consumerId, epoch,
L
Liu Jicong 已提交
1457
                 pVgCur->vgId, vgKey, buf);
H
Haojun Liao 已提交
1458
        taosHashPut(pVgOffsetHashMap, vgKey, strlen(vgKey), &pVgCur->currentOffset, sizeof(STqOffsetVal));
1459 1460 1461 1462 1463 1464 1465
      }
    }
  }

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

H
Haojun Liao 已提交
1470 1471 1472
  taosHashCleanup(pVgOffsetHashMap);

  taosThreadMutexLock(&tmq->lock);
1473
  // destroy current buffered existed topics info
1474
  if (tmq->clientTopics) {
H
Haojun Liao 已提交
1475
    taosArrayDestroyEx(tmq->clientTopics, freeClientVgInfo);
X
Xiaoyu Wang 已提交
1476
  }
1477

H
Haojun Liao 已提交
1478 1479
  tmq->clientTopics = newTopics;
  taosThreadMutexUnlock(&tmq->lock);
1480

H
Haojun Liao 已提交
1481 1482
  int8_t flag = (topicNumGet == 0)? TMQ_CONSUMER_STATUS__NO_TOPIC:TMQ_CONSUMER_STATUS__READY;
  atomic_store_8(&tmq->status, flag);
X
Xiaoyu Wang 已提交
1483
  atomic_store_32(&tmq->epoch, epoch);
H
Haojun Liao 已提交
1484

1485
  tscDebug("consumer:0x%" PRIx64 " update topic info completed", tmq->consumerId);
X
Xiaoyu Wang 已提交
1486 1487 1488
  return set;
}

H
Haojun Liao 已提交
1489
static int32_t tmqAskEpCb(void* param, SDataBuf* pMsg, int32_t code) {
1490
  SMqAskEpCbParam* pParam = (SMqAskEpCbParam*)param;
L
Liu Jicong 已提交
1491
  int8_t           async = pParam->async;
1492 1493 1494 1495 1496 1497 1498 1499 1500
  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 已提交
1501
    taosMemoryFree(pMsg->pEpSet);
1502 1503 1504 1505
    terrno = TSDB_CODE_TMQ_CONSUMER_CLOSED;
    return -1;
  }

L
Liu Jicong 已提交
1506
  pParam->code = code;
H
Haojun Liao 已提交
1507
  if (code != TSDB_CODE_SUCCESS) {
X
Xiaoyu Wang 已提交
1508 1509
    tscError("consumer:0x%" PRIx64 ", get topic endpoint error, async:%d, code:%s", tmq->consumerId, pParam->async,
             tstrerror(code));
L
Liu Jicong 已提交
1510
    goto END;
1511
  }
L
Liu Jicong 已提交
1512

L
Liu Jicong 已提交
1513
  // tmq's epoch is monotonically increase,
L
Liu Jicong 已提交
1514
  // so it's safe to discard any old epoch msg.
L
Liu Jicong 已提交
1515
  // Epoch will only increase when received newer epoch ep msg
L
Liu Jicong 已提交
1516 1517 1518
  SMqRspHead* head = pMsg->pData;
  int32_t     epoch = atomic_load_32(&tmq->epoch);
  if (head->epoch <= epoch) {
1519 1520
    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 已提交
1521
    goto END;
1522
  }
L
Liu Jicong 已提交
1523

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

L
Liu Jicong 已提交
1527
  if (!async) {
L
Liu Jicong 已提交
1528 1529
    SMqAskEpRsp rsp;
    tDecodeSMqAskEpRsp(POINTER_SHIFT(pMsg->pData, sizeof(SMqRspHead)), &rsp);
L
Liu Jicong 已提交
1530
    tmqUpdateEp(tmq, head->epoch, &rsp);
L
Liu Jicong 已提交
1531
    tDeleteSMqAskEpRsp(&rsp);
X
Xiaoyu Wang 已提交
1532
  } else {
S
Shengliang Guan 已提交
1533
    SMqAskEpRspWrapper* pWrapper = taosAllocateQitem(sizeof(SMqAskEpRspWrapper), DEF_QITEM, 0);
L
Liu Jicong 已提交
1534
    if (pWrapper == NULL) {
X
Xiaoyu Wang 已提交
1535
      terrno = TSDB_CODE_OUT_OF_MEMORY;
L
Liu Jicong 已提交
1536 1537
      code = -1;
      goto END;
X
Xiaoyu Wang 已提交
1538
    }
1539

L
Liu Jicong 已提交
1540 1541 1542
    pWrapper->tmqRspType = TMQ_MSG_TYPE__EP_RSP;
    pWrapper->epoch = head->epoch;
    memcpy(&pWrapper->msg, pMsg->pData, sizeof(SMqRspHead));
L
Liu Jicong 已提交
1543
    tDecodeSMqAskEpRsp(POINTER_SHIFT(pMsg->pData, sizeof(SMqRspHead)), &pWrapper->msg);
L
Liu Jicong 已提交
1544

L
Liu Jicong 已提交
1545
    taosWriteQitem(tmq->mqueue, pWrapper);
1546
    tsem_post(&tmq->rspSem);
1547
  }
L
Liu Jicong 已提交
1548 1549

END:
L
Liu Jicong 已提交
1550
  if (!async) {
L
Liu Jicong 已提交
1551
    tsem_post(&pParam->rspSem);
L
Liu Jicong 已提交
1552 1553
  } else {
    taosMemoryFree(pParam);
L
Liu Jicong 已提交
1554
  }
dengyihao's avatar
dengyihao 已提交
1555 1556

  taosMemoryFree(pMsg->pEpSet);
L
Liu Jicong 已提交
1557
  taosMemoryFree(pMsg->pData);
L
Liu Jicong 已提交
1558
  return code;
1559 1560
}

L
Liu Jicong 已提交
1561
void tmqBuildConsumeReqImpl(SMqPollReq* pReq, tmq_t* tmq, int64_t timeout, SMqClientTopic* pTopic, SMqClientVg* pVg) {
L
Liu Jicong 已提交
1562 1563 1564 1565
  int32_t groupLen = strlen(tmq->groupId);
  memcpy(pReq->subKey, tmq->groupId, groupLen);
  pReq->subKey[groupLen] = TMQ_SEPARATOR;
  strcpy(pReq->subKey + groupLen + 1, pTopic->topicName);
1566

1567
  pReq->withTbName = tmq->withTbName;
L
Liu Jicong 已提交
1568
  pReq->consumerId = tmq->consumerId;
1569
  pReq->timeout = timeout;
X
Xiaoyu Wang 已提交
1570
  pReq->epoch = tmq->epoch;
L
Liu Jicong 已提交
1571
  /*pReq->currentOffset = reqOffset;*/
L
Liu Jicong 已提交
1572
  pReq->reqOffset = pVg->currentOffset;
D
dapan1121 已提交
1573
  pReq->head.vgId = pVg->vgId;
1574 1575
  pReq->useSnapshot = tmq->useSnapshot;
  pReq->reqId = generateRequestId();
1576 1577
}

L
Liu Jicong 已提交
1578 1579
SMqMetaRspObj* tmqBuildMetaRspFromWrapper(SMqPollRspWrapper* pWrapper) {
  SMqMetaRspObj* pRspObj = taosMemoryCalloc(1, sizeof(SMqMetaRspObj));
L
Liu Jicong 已提交
1580
  pRspObj->resType = RES_TYPE__TMQ_META;
L
Liu Jicong 已提交
1581 1582 1583 1584 1585 1586 1587 1588
  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 已提交
1589 1590 1591
SMqRspObj* tmqBuildRspFromWrapper(SMqPollRspWrapper* pWrapper) {
  SMqRspObj* pRspObj = taosMemoryCalloc(1, sizeof(SMqRspObj));
  pRspObj->resType = RES_TYPE__TMQ;
L
Liu Jicong 已提交
1592 1593
  tstrncpy(pRspObj->topic, pWrapper->topicHandle->topicName, TSDB_TOPIC_FNAME_LEN);
  tstrncpy(pRspObj->db, pWrapper->topicHandle->db, TSDB_DB_FNAME_LEN);
L
Liu Jicong 已提交
1594
  pRspObj->vgId = pWrapper->vgHandle->vgId;
L
Liu Jicong 已提交
1595
  pRspObj->resIter = -1;
L
Liu Jicong 已提交
1596
  memcpy(&pRspObj->rsp, &pWrapper->dataRsp, sizeof(SMqDataRsp));
L
Liu Jicong 已提交
1597

L
Liu Jicong 已提交
1598 1599
  pRspObj->resInfo.totalRows = 0;
  pRspObj->resInfo.precision = TSDB_TIME_PRECISION_MILLI;
L
Liu Jicong 已提交
1600
  if (!pWrapper->dataRsp.withSchema) {
L
Liu Jicong 已提交
1601 1602
    setResSchemaInfo(&pRspObj->resInfo, pWrapper->topicHandle->schema.pSchema, pWrapper->topicHandle->schema.nCols);
  }
L
Liu Jicong 已提交
1603

L
Liu Jicong 已提交
1604
  return pRspObj;
X
Xiaoyu Wang 已提交
1605 1606
}

L
Liu Jicong 已提交
1607 1608
SMqTaosxRspObj* tmqBuildTaosxRspFromWrapper(SMqPollRspWrapper* pWrapper) {
  SMqTaosxRspObj* pRspObj = taosMemoryCalloc(1, sizeof(SMqTaosxRspObj));
1609
  pRspObj->resType = RES_TYPE__TMQ_METADATA;
L
Liu Jicong 已提交
1610 1611 1612 1613
  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;
1614
  memcpy(&pRspObj->rsp, &pWrapper->taosxRsp, sizeof(STaosxRsp));
L
Liu Jicong 已提交
1615 1616 1617

  pRspObj->resInfo.totalRows = 0;
  pRspObj->resInfo.precision = TSDB_TIME_PRECISION_MILLI;
1618
  if (!pWrapper->taosxRsp.withSchema) {
L
Liu Jicong 已提交
1619 1620 1621 1622 1623 1624
    setResSchemaInfo(&pRspObj->resInfo, pWrapper->topicHandle->schema.pSchema, pWrapper->topicHandle->schema.nCols);
  }

  return pRspObj;
}

1625 1626 1627 1628 1629 1630 1631 1632 1633 1634 1635 1636 1637 1638 1639 1640 1641 1642 1643 1644 1645 1646 1647 1648 1649 1650 1651 1652 1653 1654 1655 1656 1657
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 已提交
1658
  pParam->pVg = pVg;  // pVg may be released,fix it
1659 1660
  pParam->pTopic = pTopic;
  pParam->vgId = pVg->vgId;
H
Haojun Liao 已提交
1661
  pParam->requestId = req.reqId;
1662 1663 1664 1665 1666 1667 1668 1669 1670 1671 1672 1673 1674 1675 1676 1677 1678 1679 1680 1681 1682 1683 1684 1685

  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 已提交
1686
  tscDebug("consumer:0x%" PRIx64 " send poll to %s vgId:%d, epoch %d, req:%s, reqId:0x%" PRIx64,
1687 1688 1689 1690 1691 1692 1693 1694 1695
           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;
}

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

  for (int i = 0; i < numOfTopics; i++) {
X
Xiaoyu Wang 已提交
1702
    SMqClientTopic* pTopic = taosArrayGet(tmq->clientTopics, i);
X
Xiaoyu Wang 已提交
1703
    int32_t         numOfVg = taosArrayGetSize(pTopic->vgs);
1704 1705

    for (int j = 0; j < numOfVg; j++) {
X
Xiaoyu Wang 已提交
1706
      SMqClientVg* pVg = taosArrayGet(pTopic->vgs, j);
H
Haojun Liao 已提交
1707 1708 1709 1710 1711 1712 1713
      if (pVg->vgIgnoreCnt > 0) {
        pVg->vgIgnoreCnt -= 1;
        tscTrace("consumer:0x%" PRIx64 " epoch %d, vgId:%d idle for %d tick before poll", tmq->consumerId, tmq->epoch,
                 pVg->vgId, pVg->vgIgnoreCnt);
        continue;
      }

X
Xiaoyu Wang 已提交
1714
      int32_t      vgStatus = atomic_val_compare_exchange_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE, TMQ_VG_STATUS__WAIT);
1715
      if (vgStatus == TMQ_VG_STATUS__WAIT) {
L
Liu Jicong 已提交
1716
        int32_t vgSkipCnt = atomic_add_fetch_32(&pVg->vgSkipCnt, 1);
H
Haojun Liao 已提交
1717
        tscTrace("consumer:0x%" PRIx64 " epoch %d wait poll-rsp, skip vgId:%d skip cnt %d", tmq->consumerId, tmq->epoch,
X
Xiaoyu Wang 已提交
1718
                 pVg->vgId, vgSkipCnt);
X
Xiaoyu Wang 已提交
1719
        continue;
L
temp  
Liu Jicong 已提交
1720 1721 1722 1723
#if 0
        if (skipCnt < 30000) {
          continue;
        } else {
1724
        tscDebug("consumer:0x%" PRIx64 ",skip vgId:%d skip too much reset", tmq->consumerId, pVg->vgId);
L
temp  
Liu Jicong 已提交
1725 1726
        }
#endif
X
Xiaoyu Wang 已提交
1727
      }
1728

L
Liu Jicong 已提交
1729
      atomic_store_32(&pVg->vgSkipCnt, 0);
1730 1731 1732
      int32_t code = doTmqPollImpl(tmq, pTopic, pVg, timeout);
      if (code != TSDB_CODE_SUCCESS) {
        return code;
D
dapan1121 已提交
1733
      }
X
Xiaoyu Wang 已提交
1734 1735
    }
  }
1736

X
Xiaoyu Wang 已提交
1737 1738 1739
  return 0;
}

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

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

X
Xiaoyu Wang 已提交
1763
  while (1) {
L
Liu Jicong 已提交
1764 1765
    SMqRspWrapper* rspWrapper = NULL;
    taosGetQitem(tmq->qall, (void**)&rspWrapper);
1766

L
Liu Jicong 已提交
1767
    if (rspWrapper == NULL) {
L
Liu Jicong 已提交
1768
      taosReadAllQitems(tmq->mqueue, tmq->qall);
L
Liu Jicong 已提交
1769
      taosGetQitem(tmq->qall, (void**)&rspWrapper);
L
Liu Jicong 已提交
1770 1771 1772 1773

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

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

L
Liu Jicong 已提交
1778 1779 1780
    if (rspWrapper->tmqRspType == TMQ_MSG_TYPE__END_RSP) {
      taosFreeQitem(rspWrapper);
      terrno = TSDB_CODE_TQ_NO_COMMITTED_OFFSET;
H
Haojun Liao 已提交
1781
      tscError("consumer:0x%" PRIx64 " unexpected rsp from poll, code:%s", tmq->consumerId, tstrerror(terrno));
L
Liu Jicong 已提交
1782 1783
      return NULL;
    } else if (rspWrapper->tmqRspType == TMQ_MSG_TYPE__POLL_RSP) {
L
Liu Jicong 已提交
1784
      SMqPollRspWrapper* pollRspWrapper = (SMqPollRspWrapper*)rspWrapper;
H
Haojun Liao 已提交
1785

L
Liu Jicong 已提交
1786
      /*atomic_sub_fetch_32(&tmq->readyRequest, 1);*/
1787
      int32_t consumerEpoch = atomic_load_32(&tmq->epoch);
L
Liu Jicong 已提交
1788
      if (pollRspWrapper->dataRsp.head.epoch == consumerEpoch) {
L
Liu Jicong 已提交
1789
        SMqClientVg* pVg = pollRspWrapper->vgHandle;
L
Liu Jicong 已提交
1790
        pVg->currentOffset = pollRspWrapper->dataRsp.rspOffset;
X
Xiaoyu Wang 已提交
1791
        atomic_store_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE);
H
Haojun Liao 已提交
1792

L
Liu Jicong 已提交
1793
        if (pollRspWrapper->dataRsp.blockNum == 0) {
H
Haojun Liao 已提交
1794 1795
          tscDebug("consumer:0x%" PRIx64 " empty block received, vgId:%d, reqId:0x%" PRIx64, tmq->consumerId, pVg->vgId,
                   pollRspWrapper->reqId);
L
Liu Jicong 已提交
1796 1797
          taosFreeQitem(pollRspWrapper);
          rspWrapper = NULL;
L
temp  
Liu Jicong 已提交
1798 1799
          continue;
        }
H
Haojun Liao 已提交
1800

L
Liu Jicong 已提交
1801
        // build rsp
H
Haojun Liao 已提交
1802 1803
        char buf[80];
        tFormatOffset(buf, 80, &pVg->currentOffset);
L
Liu Jicong 已提交
1804
        SMqRspObj* pRsp = tmqBuildRspFromWrapper(pollRspWrapper);
H
Haojun Liao 已提交
1805 1806
        tscDebug("consumer:0x%" PRIx64 " process poll rsp, vgId:%d, offset:%s, blocks:%d, reqId:0x%"PRIx64, tmq->consumerId,
                     pVg->vgId, buf, pollRspWrapper->dataRsp.blockNum, pollRspWrapper->reqId);
H
Haojun Liao 已提交
1807

L
Liu Jicong 已提交
1808
        taosFreeQitem(pollRspWrapper);
L
Liu Jicong 已提交
1809
        return pRsp;
X
Xiaoyu Wang 已提交
1810
      } else {
X
Xiaoyu Wang 已提交
1811
        tscDebug("consumer:0x%" PRIx64 " msg discard since epoch mismatch: msg epoch %d, consumer epoch %d",
1812
                 tmq->consumerId, pollRspWrapper->dataRsp.head.epoch, consumerEpoch);
L
Liu Jicong 已提交
1813
        tmqFreeRspWrapper(rspWrapper);
L
Liu Jicong 已提交
1814 1815 1816 1817 1818
        taosFreeQitem(pollRspWrapper);
      }
    } else if (rspWrapper->tmqRspType == TMQ_MSG_TYPE__POLL_META_RSP) {
      SMqPollRspWrapper* pollRspWrapper = (SMqPollRspWrapper*)rspWrapper;
      int32_t            consumerEpoch = atomic_load_32(&tmq->epoch);
1819 1820 1821

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

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

L
Liu Jicong 已提交
1840 1841 1842 1843
      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 已提交
1844

L
Liu Jicong 已提交
1845 1846
        if (pollRspWrapper->taosxRsp.blockNum == 0) {
          rspWrapper = NULL;
H
Haojun Liao 已提交
1847 1848
          tscDebug("consumer:0x%" PRIx64 " taosx empty block received, vgId:%d, reqId:0x%" PRIx64, tmq->consumerId, pVg->vgId,
                   pollRspWrapper->reqId);
H
Haojun Liao 已提交
1849
          pVg->vgIgnoreCnt = VG_POLL_IGNORE_TICK;
H
Haojun Liao 已提交
1850
          taosFreeQitem(pollRspWrapper);
L
Liu Jicong 已提交
1851 1852
          continue;
        }
wmmhello's avatar
wmmhello 已提交
1853

L
Liu Jicong 已提交
1854
        // build rsp
wmmhello's avatar
wmmhello 已提交
1855
        void* pRsp = NULL;
L
Liu Jicong 已提交
1856
        if (pollRspWrapper->taosxRsp.createTableNum == 0) {
wmmhello's avatar
wmmhello 已提交
1857
          pRsp = tmqBuildRspFromWrapper(pollRspWrapper);
L
Liu Jicong 已提交
1858
        } else {
wmmhello's avatar
wmmhello 已提交
1859 1860
          pRsp = tmqBuildTaosxRspFromWrapper(pollRspWrapper);
        }
H
Haojun Liao 已提交
1861 1862 1863 1864 1865 1866


        char buf[80];
        tFormatOffset(buf, 80, &pVg->currentOffset);
        tscDebug("consumer:0x%" PRIx64 " process taosx poll rsp, vgId:%d, offset:%s, blocks:%d, reqId:0x%"PRIx64, tmq->consumerId, pVg->vgId,
                 buf, pollRspWrapper->dataRsp.blockNum, pollRspWrapper->reqId);
H
Haojun Liao 已提交
1867 1868

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

L
Liu Jicong 已提交
1871
      } else {
X
Xiaoyu Wang 已提交
1872
        tscDebug("consumer:0x%" PRIx64 " msg discard since epoch mismatch: msg epoch %d, consumer epoch %d",
1873
                 tmq->consumerId, pollRspWrapper->taosxRsp.head.epoch, consumerEpoch);
L
Liu Jicong 已提交
1874
        tmqFreeRspWrapper(rspWrapper);
L
Liu Jicong 已提交
1875 1876
        taosFreeQitem(pollRspWrapper);
      }
X
Xiaoyu Wang 已提交
1877
    } else {
H
Haojun Liao 已提交
1878 1879
      tscDebug("consumer:0x%" PRIx64 " not data msg received", tmq->consumerId);

X
Xiaoyu Wang 已提交
1880
      bool reset = false;
L
Liu Jicong 已提交
1881 1882
      tmqHandleNoPollRsp(tmq, rspWrapper, &reset);
      taosFreeQitem(rspWrapper);
X
Xiaoyu Wang 已提交
1883
      if (pollIfReset && reset) {
1884
        tscDebug("consumer:0x%" PRIx64 ", reset and repoll", tmq->consumerId);
1885
        tmqPollImpl(tmq, timeout);
X
Xiaoyu Wang 已提交
1886 1887 1888 1889 1890
      }
    }
  }
}

1891
TAOS_RES* tmq_consumer_poll(tmq_t* tmq, int64_t timeout) {
L
Liu Jicong 已提交
1892 1893
  void*   rspObj;
  int64_t startTime = taosGetTimestampMs();
L
Liu Jicong 已提交
1894

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

1897 1898 1899
#if 0
  tmqHandleAllDelayedTask(tmq);
  tmqPollImpl(tmq, timeout);
1900
  rspObj = tmqHandleAllRsp(tmq, timeout, false);
L
Liu Jicong 已提交
1901 1902
  if (rspObj) {
    return (TAOS_RES*)rspObj;
L
fix  
Liu Jicong 已提交
1903
  }
1904
#endif
X
Xiaoyu Wang 已提交
1905

1906
  // in no topic status, delayed task also need to be processed
L
Liu Jicong 已提交
1907
  if (atomic_load_8(&tmq->status) == TMQ_CONSUMER_STATUS__INIT) {
1908
    tscDebug("consumer:0x%" PRIx64 " poll return since consumer is init", tmq->consumerId);
1909
    taosMsleep(500);  //     sleep for a while
1910 1911 1912
    return NULL;
  }

L
Liu Jicong 已提交
1913 1914 1915
  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 已提交
1916
      if (retryCnt++ > 40) {
L
Liu Jicong 已提交
1917 1918
        return NULL;
      }
1919

H
Haojun Liao 已提交
1920
      tscDebug("consumer:0x%" PRIx64 " not ready, retry:%d/40 in 500ms", tmq->consumerId, retryCnt);
L
Liu Jicong 已提交
1921 1922 1923 1924
      taosMsleep(500);
    }
  }

X
Xiaoyu Wang 已提交
1925
  while (1) {
L
Liu Jicong 已提交
1926
    tmqHandleAllDelayedTask(tmq);
1927

L
Liu Jicong 已提交
1928
    if (tmqPollImpl(tmq, timeout) < 0) {
1929
      tscDebug("consumer:0x%" PRIx64 " return due to poll error", tmq->consumerId);
L
Liu Jicong 已提交
1930 1931
      /*return NULL;*/
    }
L
Liu Jicong 已提交
1932

1933
    rspObj = tmqHandleAllRsp(tmq, timeout, false);
L
Liu Jicong 已提交
1934
    if (rspObj) {
1935
      tscDebug("consumer:0x%" PRIx64 " return rsp %p", tmq->consumerId, rspObj);
L
Liu Jicong 已提交
1936
      return (TAOS_RES*)rspObj;
L
Liu Jicong 已提交
1937
    } else if (terrno == TSDB_CODE_TQ_NO_COMMITTED_OFFSET) {
1938
      tscDebug("consumer:0x%" PRIx64 " return null since no committed offset", tmq->consumerId);
L
Liu Jicong 已提交
1939
      return NULL;
X
Xiaoyu Wang 已提交
1940
    }
1941

1942
    if (timeout != -1) {
L
Liu Jicong 已提交
1943
      int64_t currentTime = taosGetTimestampMs();
1944 1945 1946
      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 已提交
1947
                 tmq->consumerId, tmq->epoch, startTime, currentTime);
X
Xiaoyu Wang 已提交
1948 1949
        return NULL;
      }
1950
      /*tscInfo("consumer:0x%" PRIx64 ", (epoch %d) wait, start time %" PRId64 ", current time %" PRId64*/
L
Liu Jicong 已提交
1951
      /*", left time %" PRId64,*/
1952 1953
      /*tmq->consumerId, tmq->epoch, startTime, currentTime, (timeout - elapsedTime));*/
      tsem_timewait(&tmq->rspSem, (timeout - elapsedTime));
L
Liu Jicong 已提交
1954 1955
    } else {
      // use tsem_timewait instead of tsem_wait to avoid unexpected stuck
L
Liu Jicong 已提交
1956
      tsem_timewait(&tmq->rspSem, 1000);
X
Xiaoyu Wang 已提交
1957 1958 1959 1960
    }
  }
}

L
Liu Jicong 已提交
1961
int32_t tmq_consumer_close(tmq_t* tmq) {
1962
  if (tmq->status == TMQ_CONSUMER_STATUS__READY) {
L
Liu Jicong 已提交
1963 1964
    int32_t rsp = tmq_commit_sync(tmq, NULL);
    if (rsp != 0) {
L
Liu Jicong 已提交
1965
      return rsp;
1966 1967
    }

L
Liu Jicong 已提交
1968
    int32_t     retryCnt = 0;
1969
    tmq_list_t* lst = tmq_list_new();
L
Liu Jicong 已提交
1970 1971 1972 1973 1974 1975 1976 1977 1978 1979
    while (1) {
      rsp = tmq_subscribe(tmq, lst);
      if (rsp != TSDB_CODE_MND_CONSUMER_NOT_READY || retryCnt > 5) {
        break;
      } else {
        retryCnt++;
        taosMsleep(500);
      }
    }

1980
    tmq_list_destroy(lst);
L
Liu Jicong 已提交
1981
  }
H
Haojun Liao 已提交
1982

1983
  taosRemoveRef(tmqMgmt.rsetId, tmq->refId);
L
Liu Jicong 已提交
1984
  return 0;
1985
}
L
Liu Jicong 已提交
1986

L
Liu Jicong 已提交
1987 1988
const char* tmq_err2str(int32_t err) {
  if (err == 0) {
L
Liu Jicong 已提交
1989
    return "success";
L
Liu Jicong 已提交
1990
  } else if (err == -1) {
L
Liu Jicong 已提交
1991 1992 1993
    return "fail";
  } else {
    return tstrerror(err);
L
Liu Jicong 已提交
1994 1995
  }
}
L
Liu Jicong 已提交
1996

L
Liu Jicong 已提交
1997 1998 1999 2000 2001
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;
2002 2003
  } else if (TD_RES_TMQ_METADATA(res)) {
    return TMQ_RES_METADATA;
L
Liu Jicong 已提交
2004 2005 2006 2007 2008
  } else {
    return TMQ_RES_INVALID;
  }
}

L
Liu Jicong 已提交
2009
const char* tmq_get_topic_name(TAOS_RES* res) {
L
Liu Jicong 已提交
2010 2011
  if (TD_RES_TMQ(res)) {
    SMqRspObj* pRspObj = (SMqRspObj*)res;
L
Liu Jicong 已提交
2012
    return strchr(pRspObj->topic, '.') + 1;
L
Liu Jicong 已提交
2013 2014 2015
  } else if (TD_RES_TMQ_META(res)) {
    SMqMetaRspObj* pMetaRspObj = (SMqMetaRspObj*)res;
    return strchr(pMetaRspObj->topic, '.') + 1;
2016 2017 2018
  } else if (TD_RES_TMQ_METADATA(res)) {
    SMqTaosxRspObj* pRspObj = (SMqTaosxRspObj*)res;
    return strchr(pRspObj->topic, '.') + 1;
L
Liu Jicong 已提交
2019 2020 2021 2022 2023
  } else {
    return NULL;
  }
}

L
Liu Jicong 已提交
2024 2025 2026 2027
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 已提交
2028 2029 2030
  } else if (TD_RES_TMQ_META(res)) {
    SMqMetaRspObj* pMetaRspObj = (SMqMetaRspObj*)res;
    return strchr(pMetaRspObj->db, '.') + 1;
2031 2032 2033
  } else if (TD_RES_TMQ_METADATA(res)) {
    SMqTaosxRspObj* pRspObj = (SMqTaosxRspObj*)res;
    return strchr(pRspObj->db, '.') + 1;
L
Liu Jicong 已提交
2034 2035 2036 2037 2038
  } else {
    return NULL;
  }
}

L
Liu Jicong 已提交
2039 2040 2041 2042
int32_t tmq_get_vgroup_id(TAOS_RES* res) {
  if (TD_RES_TMQ(res)) {
    SMqRspObj* pRspObj = (SMqRspObj*)res;
    return pRspObj->vgId;
L
Liu Jicong 已提交
2043 2044 2045
  } else if (TD_RES_TMQ_META(res)) {
    SMqMetaRspObj* pMetaRspObj = (SMqMetaRspObj*)res;
    return pMetaRspObj->vgId;
2046
  } else if (TD_RES_TMQ_METADATA(res)) {
2047 2048
    SMqTaosxRspObj* pRspObj = (SMqTaosxRspObj*)res;
    return pRspObj->vgId;
L
Liu Jicong 已提交
2049 2050 2051 2052
  } else {
    return -1;
  }
}
L
Liu Jicong 已提交
2053 2054 2055 2056 2057 2058 2059 2060

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;
    }
2061
    return (const char*)taosArrayGetP(pRspObj->rsp.blockTbName, pRspObj->resIter);
2062 2063
  } else if (TD_RES_TMQ_METADATA(res)) {
    SMqTaosxRspObj* pRspObj = (SMqTaosxRspObj*)res;
L
Liu Jicong 已提交
2064 2065 2066
    if (!pRspObj->rsp.withTbName || pRspObj->rsp.blockTbName == NULL || pRspObj->resIter < 0 ||
        pRspObj->resIter >= pRspObj->rsp.blockNum) {
      return NULL;
2067
    }
L
Liu Jicong 已提交
2068 2069
    return (const char*)taosArrayGetP(pRspObj->rsp.blockTbName, pRspObj->resIter);
  }
L
Liu Jicong 已提交
2070 2071
  return NULL;
}
2072

L
Liu Jicong 已提交
2073
void tmq_commit_async(tmq_t* tmq, const TAOS_RES* msg, tmq_commit_cb* cb, void* param) {
L
Liu Jicong 已提交
2074
  tmqCommitInner(tmq, msg, 0, 1, cb, param);
L
Liu Jicong 已提交
2075 2076
}

2077
int32_t tmq_commit_sync(tmq_t* tmq, const TAOS_RES* msg) {
L
Liu Jicong 已提交
2078
  return tmqCommitInner(tmq, msg, 0, 0, NULL, NULL);
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 2115 2116 2117 2118 2119 2120 2121 2122 2123 2124 2125 2126 2127 2128 2129 2130 2131 2132 2133 2134 2135 2136 2137 2138 2139 2140 2141 2142 2143 2144 2145 2146 2147 2148 2149 2150 2151 2152 2153 2154 2155 2156 2157 2158 2159 2160 2161 2162 2163 2164 2165 2166 2167 2168 2169 2170 2171 2172 2173 2174 2175 2176 2177 2178 2179 2180 2181 2182 2183 2184

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);
2185
    } else if (!pParamSet->automatic && pParamSet->userCb) { // sem post
2186 2187
      pParamSet->userCb(tmq, pParamSet->rspErr, pParamSet->userParam);
    }
2188

2189 2190 2191 2192 2193 2194 2195 2196 2197 2198
    taosMemoryFree(pParamSet);
  } else {
    tsem_post(&pParamSet->rspSem);
  }

#if 0
  taosArrayDestroyP(pParamSet->successfulOffsets, taosMemoryFree);
    taosArrayDestroyP(pParamSet->failedOffsets, taosMemoryFree);
#endif
  return 0;
2199 2200 2201 2202 2203
}

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 已提交
2204 2205
    tscDebug("consumer:0x%" PRIx64 " topic:%s vgId:%d all commit-rsp received, commit completed", consumerId, pTopic,
             vgId);
2206
    tmqCommitDone(pParamSet);
H
Haojun Liao 已提交
2207 2208 2209
  } else {
    tscDebug("consumer:0x%" PRIx64 " topic:%s vgId:%d commit-rsp received, remain:%d", consumerId, pTopic, vgId,
             waitingRspNum);
2210 2211
  }
}