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

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

X
Xiaoyu Wang 已提交
27 28
#define EMPTY_BLOCK_POLL_IDLE_DURATION 10
#define DEFAULT_AUTO_COMMIT_INTERVAL   5000
29

30 31
#define OFFSET_IS_RESET_OFFSET(_of)  ((_of) < 0)

32 33
typedef void (*__tmq_askep_fn_t)(tmq_t* pTmq, int32_t code, SDataBuf* pBuf, void* pParam);

X
Xiaoyu Wang 已提交
34
struct SMqMgmt {
35 36 37
  int8_t  inited;
  tmr_h   timer;
  int32_t rsetId;
38
};
L
Liu Jicong 已提交
39

X
Xiaoyu Wang 已提交
40 41
static TdThreadOnce   tmqInit = PTHREAD_ONCE_INIT;  // initialize only once
volatile int32_t      tmqInitRes = 0;               // initialize rsp code
42
static struct SMqMgmt tmqMgmt = {0};
43

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

typedef struct {
L
Liu Jicong 已提交
50 51 52
  int8_t      tmqRspType;
  int32_t     epoch;
  SMqAskEpRsp msg;
L
Liu Jicong 已提交
53 54
} SMqAskEpRspWrapper;

L
Liu Jicong 已提交
55
struct tmq_list_t {
L
Liu Jicong 已提交
56
  SArray container;
L
Liu Jicong 已提交
57
};
L
Liu Jicong 已提交
58

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

struct tmq_t {
78 79 80 81 82 83 84
  int64_t        refId;
  char           groupId[TSDB_CGROUP_LEN];
  char           clientId[256];
  int8_t         withTbName;
  int8_t         useSnapshot;
  int8_t         autoCommit;
  int32_t        autoCommitInterval;
85
  int8_t         resetOffsetCfg;
86 87
  uint64_t       consumerId;
  bool           hbBgEnable;
L
Liu Jicong 已提交
88 89
  tmq_commit_cb* commitCb;
  void*          commitCbUserParam;
L
Liu Jicong 已提交
90 91

  // status
wmmhello's avatar
wmmhello 已提交
92
  SRWLatch        lock;
L
Liu Jicong 已提交
93 94
  int8_t  status;
  int32_t epoch;
L
Liu Jicong 已提交
95 96
#if 0
  int8_t  epStatus;
L
Liu Jicong 已提交
97
  int32_t epSkipCnt;
L
Liu Jicong 已提交
98
#endif
99
  // poll info
X
Xiaoyu Wang 已提交
100 101
  int64_t pollCnt;
  int64_t totalRows;
wmmhello's avatar
wmmhello 已提交
102
//  bool    needReportOffsetRows;
L
Liu Jicong 已提交
103

L
Liu Jicong 已提交
104
  // timer
X
Xiaoyu Wang 已提交
105 106 107 108 109 110 111 112 113 114
  tmr_h       hbLiveTimer;
  tmr_h       epTimer;
  tmr_h       reportTimer;
  tmr_h       commitTimer;
  STscObj*    pTscObj;       // connection
  SArray*     clientTopics;  // SArray<SMqClientTopic>
  STaosQueue* mqueue;        // queue of rsp
  STaosQall*  qall;
  STaosQueue* delayedTask;  // delayed task queue for heartbeat and auto commit
  tsem_t      rspSem;
L
Liu Jicong 已提交
115 116
};

117 118
typedef struct SAskEpInfo {
  int32_t code;
H
Haojun Liao 已提交
119
  tsem_t  sem;
120 121
} SAskEpInfo;

X
Xiaoyu Wang 已提交
122 123 124 125 126 127 128 129
enum {
  TMQ_VG_STATUS__IDLE = 0,
  TMQ_VG_STATUS__WAIT,
};

enum {
  TMQ_CONSUMER_STATUS__INIT = 0,
  TMQ_CONSUMER_STATUS__READY,
130
  TMQ_CONSUMER_STATUS__NO_TOPIC,
L
Liu Jicong 已提交
131
  TMQ_CONSUMER_STATUS__RECOVER,
L
Liu Jicong 已提交
132 133
};

L
Liu Jicong 已提交
134
enum {
135
  TMQ_DELAYED_TASK__ASK_EP = 1,
L
Liu Jicong 已提交
136 137 138 139
  TMQ_DELAYED_TASK__REPORT,
  TMQ_DELAYED_TASK__COMMIT,
};

H
Haojun Liao 已提交
140
typedef struct SVgOffsetInfo {
L
Liu Jicong 已提交
141 142
  STqOffsetVal committedOffset;
  STqOffsetVal currentOffset;
H
Haojun Liao 已提交
143 144 145 146 147 148 149 150 151 152
  int64_t      walVerBegin;
  int64_t      walVerEnd;
} SVgOffsetInfo;

typedef struct {
  int64_t       pollCnt;
  int64_t       numOfRows;
  SVgOffsetInfo offsetInfo;
  int32_t       vgId;
  int32_t       vgStatus;
H
Haojun Liao 已提交
153 154 155 156
  int32_t       vgSkipCnt;              // here used to mark the slow vgroups
  bool          receivedInfoFromVnode;  // has already received info from vnode
  int64_t       emptyBlockReceiveTs;    // once empty block is received, idle for ignoreCnt then start to poll data
  bool          seekUpdated;            // offset is updated by seek operator, therefore, not update by vnode rsp.
H
Haojun Liao 已提交
157
  SEpSet        epSet;
158 159
} SMqClientVg;

L
Liu Jicong 已提交
160
typedef struct {
161 162 163
  char           topicName[TSDB_TOPIC_FNAME_LEN];
  char           db[TSDB_DB_FNAME_LEN];
  SArray*        vgs;  // SArray<SMqClientVg>
L
Liu Jicong 已提交
164
  SSchemaWrapper schema;
165 166
} SMqClientTopic;

L
Liu Jicong 已提交
167 168
typedef struct {
  int8_t          tmqRspType;
169
  int32_t         epoch;  // epoch can be used to guard the vgHandle
170
  int32_t         vgId;
wmmhello's avatar
wmmhello 已提交
171
  char            topicName[TSDB_TOPIC_FNAME_LEN];
L
Liu Jicong 已提交
172 173
  SMqClientVg*    vgHandle;
  SMqClientTopic* topicHandle;
H
Haojun Liao 已提交
174
  uint64_t        reqId;
175
  SEpSet*         pEpset;
L
Liu Jicong 已提交
176
  union {
L
Liu Jicong 已提交
177 178
    SMqDataRsp dataRsp;
    SMqMetaRsp metaRsp;
L
Liu Jicong 已提交
179
    STaosxRsp  taosxRsp;
L
Liu Jicong 已提交
180
  };
L
Liu Jicong 已提交
181 182
} SMqPollRspWrapper;

L
Liu Jicong 已提交
183
typedef struct {
wmmhello's avatar
wmmhello 已提交
184 185
//  int64_t refId;
//  int32_t epoch;
L
Liu Jicong 已提交
186 187
  tsem_t  rspSem;
  int32_t rspErr;
L
Liu Jicong 已提交
188
} SMqSubscribeCbParam;
L
Liu Jicong 已提交
189

L
Liu Jicong 已提交
190
typedef struct {
191 192 193 194
  int64_t          refId;
  int32_t          epoch;
  void*            pParam;
  __tmq_askep_fn_t pUserFn;
195 196
} SMqAskEpCbParam;

L
Liu Jicong 已提交
197
typedef struct {
198 199
  int64_t         refId;
  int32_t         epoch;
wmmhello's avatar
wmmhello 已提交
200 201 202
  char            topicName[TSDB_TOPIC_FNAME_LEN];
//  SMqClientVg*    pVg;
//  SMqClientTopic* pTopic;
L
Liu Jicong 已提交
203
  int32_t         vgId;
X
Xiaoyu Wang 已提交
204
  uint64_t        requestId;  // request id for debug purpose
X
Xiaoyu Wang 已提交
205
} SMqPollCbParam;
206

207 208 209 210 211 212 213 214 215 216 217 218 219 220 221 222 223
typedef struct SMqVgCommon {
  tsem_t        rsp;
  int32_t       numOfRsp;
  SArray*       pList;
  TdThreadMutex mutex;
  int64_t       consumerId;
  char*         pTopicName;
  int32_t       code;
} SMqVgCommon;

typedef struct SMqVgWalInfoParam {
  int32_t      vgId;
  int32_t      epoch;
  int32_t      totalReq;
  SMqVgCommon* pCommon;
} SMqVgWalInfoParam;

224
typedef struct {
225 226
  int64_t        refId;
  int32_t        epoch;
L
Liu Jicong 已提交
227 228
  int32_t        waitingRspNum;
  int32_t        totalRspNum;
229
  int32_t        code;
230
  tmq_commit_cb* callbackFn;
L
Liu Jicong 已提交
231 232
  /*SArray*        successfulOffsets;*/
  /*SArray*        failedOffsets;*/
X
Xiaoyu Wang 已提交
233
  void* userParam;
234 235 236 237
} SMqCommitCbParamSet;

typedef struct {
  SMqCommitCbParamSet* params;
238
  SMqVgOffset*         pOffset;
H
Haojun Liao 已提交
239 240 241
  char                 topicName[TSDB_TOPIC_FNAME_LEN];
  int32_t              vgId;
  tmq_t*               pTmq;
242
} SMqCommitCbParam;
243

244 245 246 247 248
typedef struct SSyncCommitInfo {
  tsem_t  sem;
  int32_t code;
} SSyncCommitInfo;

249
static int32_t doAskEp(tmq_t* tmq);
250 251
static int32_t makeTopicVgroupKey(char* dst, const char* topicName, int32_t vg);
static int32_t tmqCommitDone(SMqCommitCbParamSet* pParamSet);
252
static int32_t doSendCommitMsg(tmq_t* tmq, SMqClientVg* pVg, const char* pTopicName, SMqCommitCbParamSet* pParamSet,
253
                               int32_t index, int32_t totalVgroups, int32_t type);
X
Xiaoyu Wang 已提交
254 255 256
static void    commitRspCountDown(SMqCommitCbParamSet* pParamSet, int64_t consumerId, const char* pTopic, int32_t vgId);
static void    asyncAskEp(tmq_t* pTmq, __tmq_askep_fn_t askEpFn, void* param);
static void    addToQueueCallbackFn(tmq_t* pTmq, int32_t code, SDataBuf* pDataBuf, void* param);
257

258
tmq_conf_t* tmq_conf_new() {
wafwerar's avatar
wafwerar 已提交
259
  tmq_conf_t* conf = taosMemoryCalloc(1, sizeof(tmq_conf_t));
260 261 262 263 264
  if (conf == NULL) {
    terrno = TSDB_CODE_OUT_OF_MEMORY;
    return conf;
  }

265
  conf->withTbName = false;
L
Liu Jicong 已提交
266
  conf->autoCommit = true;
267
  conf->autoCommitInterval = DEFAULT_AUTO_COMMIT_INTERVAL;
268
  conf->resetOffset = TMQ_OFFSET__RESET_EARLIEST;
269
  conf->hbBgEnable = true;
270

271 272 273
  return conf;
}

L
Liu Jicong 已提交
274
void tmq_conf_destroy(tmq_conf_t* conf) {
L
Liu Jicong 已提交
275
  if (conf) {
276 277 278 279 280 281 282 283 284
    if (conf->ip) {
      taosMemoryFree(conf->ip);
    }
    if (conf->user) {
      taosMemoryFree(conf->user);
    }
    if (conf->pass) {
      taosMemoryFree(conf->pass);
    }
L
Liu Jicong 已提交
285 286
    taosMemoryFree(conf);
  }
L
Liu Jicong 已提交
287 288 289
}

tmq_conf_res_t tmq_conf_set(tmq_conf_t* conf, const char* key, const char* value) {
290
  if (strcasecmp(key, "group.id") == 0) {
L
Liu Jicong 已提交
291
    tstrncpy(conf->groupId, value, TSDB_CGROUP_LEN);
L
Liu Jicong 已提交
292
    return TMQ_CONF_OK;
293
  }
L
Liu Jicong 已提交
294

295
  if (strcasecmp(key, "client.id") == 0) {
L
Liu Jicong 已提交
296
    tstrncpy(conf->clientId, value, 256);
L
Liu Jicong 已提交
297 298
    return TMQ_CONF_OK;
  }
L
Liu Jicong 已提交
299

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

312
  if (strcasecmp(key, "auto.commit.interval.ms") == 0) {
313
    conf->autoCommitInterval = taosStr2int64(value);
L
Liu Jicong 已提交
314 315 316
    return TMQ_CONF_OK;
  }

317 318 319
  if (strcasecmp(key, "auto.offset.reset") == 0) {
    if (strcasecmp(value, "none") == 0) {
      conf->resetOffset = TMQ_OFFSET__RESET_NONE;
L
Liu Jicong 已提交
320
      return TMQ_CONF_OK;
321
    } else if (strcasecmp(value, "earliest") == 0) {
322
      conf->resetOffset = TMQ_OFFSET__RESET_EARLIEST;
L
Liu Jicong 已提交
323
      return TMQ_CONF_OK;
324 325
    } else if (strcasecmp(value, "latest") == 0) {
      conf->resetOffset = TMQ_OFFSET__RESET_LATEST;
L
Liu Jicong 已提交
326 327 328 329 330
      return TMQ_CONF_OK;
    } else {
      return TMQ_CONF_INVALID;
    }
  }
L
Liu Jicong 已提交
331

332 333
  if (strcasecmp(key, "msg.with.table.name") == 0) {
    if (strcasecmp(value, "true") == 0) {
334
      conf->withTbName = true;
L
Liu Jicong 已提交
335
      return TMQ_CONF_OK;
336
    } else if (strcasecmp(value, "false") == 0) {
337
      conf->withTbName = false;
L
Liu Jicong 已提交
338
      return TMQ_CONF_OK;
339 340 341 342 343
    } else {
      return TMQ_CONF_INVALID;
    }
  }

344 345
  if (strcasecmp(key, "experimental.snapshot.enable") == 0) {
    if (strcasecmp(value, "true") == 0) {
L
Liu Jicong 已提交
346
      conf->snapEnable = true;
347
      return TMQ_CONF_OK;
348
    } else if (strcasecmp(value, "false") == 0) {
L
Liu Jicong 已提交
349
      conf->snapEnable = false;
350 351 352 353 354 355
      return TMQ_CONF_OK;
    } else {
      return TMQ_CONF_INVALID;
    }
  }

356
  if (strcasecmp(key, "experimental.snapshot.batch.size") == 0) {
357
    conf->snapBatchSize = taosStr2int64(value);
L
Liu Jicong 已提交
358 359 360
    return TMQ_CONF_OK;
  }

361
//  if (strcasecmp(key, "enable.heartbeat.background") == 0) {
X
Xiaoyu Wang 已提交
362 363 364 365 366 367 368
    //    if (strcasecmp(value, "true") == 0) {
    //      conf->hbBgEnable = true;
    //      return TMQ_CONF_OK;
    //    } else if (strcasecmp(value, "false") == 0) {
    //      conf->hbBgEnable = false;
    //      return TMQ_CONF_OK;
    //    } else {
369 370
//    tscError("the default value of enable.heartbeat.background is true, can not be seted");
//    return TMQ_CONF_INVALID;
X
Xiaoyu Wang 已提交
371
    //    }
372
//  }
L
Liu Jicong 已提交
373

374
  if (strcasecmp(key, "td.connect.ip") == 0) {
375
    conf->ip = taosStrdup(value);
L
Liu Jicong 已提交
376 377
    return TMQ_CONF_OK;
  }
378

379
  if (strcasecmp(key, "td.connect.user") == 0) {
380
    conf->user = taosStrdup(value);
L
Liu Jicong 已提交
381 382
    return TMQ_CONF_OK;
  }
383

384
  if (strcasecmp(key, "td.connect.pass") == 0) {
385
    conf->pass = taosStrdup(value);
L
Liu Jicong 已提交
386 387
    return TMQ_CONF_OK;
  }
388

389
  if (strcasecmp(key, "td.connect.port") == 0) {
390
    conf->port = taosStr2int64(value);
L
Liu Jicong 已提交
391 392
    return TMQ_CONF_OK;
  }
393

394
  if (strcasecmp(key, "td.connect.db") == 0) {
L
Liu Jicong 已提交
395 396 397
    return TMQ_CONF_OK;
  }

L
Liu Jicong 已提交
398
  return TMQ_CONF_UNKNOWN;
399 400
}

X
Xiaoyu Wang 已提交
401
tmq_list_t* tmq_list_new() { return (tmq_list_t*)taosArrayInit(0, sizeof(void*)); }
402

L
Liu Jicong 已提交
403 404
int32_t tmq_list_append(tmq_list_t* list, const char* src) {
  SArray* container = &list->container;
405
  if (src == NULL || src[0] == 0) return -1;
406
  char* topic = taosStrdup(src);
L
fix  
Liu Jicong 已提交
407
  if (taosArrayPush(container, &topic) == NULL) return -1;
408 409 410
  return 0;
}

L
Liu Jicong 已提交
411
void tmq_list_destroy(tmq_list_t* list) {
L
Liu Jicong 已提交
412
  SArray* container = &list->container;
L
Liu Jicong 已提交
413
  taosArrayDestroyP(container, taosMemoryFree);
L
Liu Jicong 已提交
414 415
}

L
Liu Jicong 已提交
416 417 418 419 420 421 422 423 424 425
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;
}

426 427 428 429 430 431 432 433 434 435 436 437 438 439 440 441 442 443 444 445 446 447 448 449
//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;
//
//  for (int32_t i = 0; i < numOfTopics; ++i) {
//    SMqClientTopic* pTopic = taosArrayGet(pTopicList, i);
//    if (strcmp(pTopic->topicName, pName) != 0) {
//      continue;
//    }
//
//    *numOfVgroups = taosArrayGetSize(pTopic->vgs);
//    for (int32_t j = 0; j < (*numOfVgroups); ++j) {
//      SMqClientVg* pClientVg = taosArrayGet(pTopic->vgs, j);
//      if (pClientVg->vgId == vgId) {
//        *index = j;
//        return pClientVg;
//      }
//    }
//  }
//
//  return NULL;
//}
450

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

X
Xiaoyu Wang 已提交
458 459 460 461 462 463 464 465 466 467 468 469 470 471
  //  if (code != TSDB_CODE_SUCCESS) { // if commit offset failed, let's try again
  //    taosThreadMutexLock(&pParam->pTmq->lock);
  //    int32_t numOfVgroups, index;
  //    SMqClientVg* pVg = foundClientVg(pParam->pTmq->clientTopics, pParam->topicName, pParam->vgId, &index,
  //    &numOfVgroups); if (pVg == NULL) {
  //      tscDebug("consumer:0x%" PRIx64
  //               " subKey:%s vgId:%d commit failed, code:%s has been transferred to other consumer, no need retry
  //               ordinal:%d/%d", pParam->pTmq->consumerId, pParam->pOffset->subKey, pParam->vgId, tstrerror(code),
  //               index + 1, numOfVgroups);
  //    } else { // let's retry the commit
  //      int32_t code1 = doSendCommitMsg(pParam->pTmq, pVg, pParam->topicName, pParamSet, index, numOfVgroups);
  //      if (code1 != TSDB_CODE_SUCCESS) {  // retry failed.
  //        tscError("consumer:0x%" PRIx64 " topic:%s vgId:%d offset:%" PRId64
  //                 " retry failed, ignore this commit. code:%s ordinal:%d/%d",
H
Haojun Liao 已提交
472
  //                 pParam->pTmq->consumerId, pParam->topicName, pVg->vgId, pVg->offsetInfo.committedOffset.version,
X
Xiaoyu Wang 已提交
473 474 475 476 477 478 479 480 481 482 483 484 485 486 487
  //                 tstrerror(terrno), index + 1, numOfVgroups);
  //      }
  //    }
  //
  //    taosThreadMutexUnlock(&pParam->pTmq->lock);
  //
  //    taosMemoryFree(pParam->pOffset);
  //    taosMemoryFree(pBuf->pData);
  //    taosMemoryFree(pBuf->pEpSet);
  //
  //    commitRspCountDown(pParamSet, pParam->pTmq->consumerId, pParam->topicName, pParam->vgId);
  //    return 0;
  //  }
  //
  //  // todo replace the pTmq with refId
488

L
Liu Jicong 已提交
489
  taosMemoryFree(pParam->pOffset);
L
Liu Jicong 已提交
490
  taosMemoryFree(pBuf->pData);
dengyihao's avatar
dengyihao 已提交
491
  taosMemoryFree(pBuf->pEpSet);
L
Liu Jicong 已提交
492

493
  commitRspCountDown(pParamSet, pParam->pTmq->consumerId, pParam->topicName, pParam->vgId);
494 495 496
  return 0;
}

497
static int32_t doSendCommitMsg(tmq_t* tmq, SMqClientVg* pVg, const char* pTopicName, SMqCommitCbParamSet* pParamSet,
498
                               int32_t index, int32_t totalVgroups, int32_t type) {
499
  SMqVgOffset* pOffset = taosMemoryCalloc(1, sizeof(SMqVgOffset));
L
Liu Jicong 已提交
500
  if (pOffset == NULL) {
501
    return TSDB_CODE_OUT_OF_MEMORY;
L
Liu Jicong 已提交
502
  }
503

504 505
  pOffset->consumerId = tmq->consumerId;
  pOffset->offset.val = pVg->offsetInfo.currentOffset;
506

L
Liu Jicong 已提交
507
  int32_t groupLen = strlen(tmq->groupId);
508 509 510
  memcpy(pOffset->offset.subKey, tmq->groupId, groupLen);
  pOffset->offset.subKey[groupLen] = TMQ_SEPARATOR;
  strcpy(pOffset->offset.subKey + groupLen + 1, pTopicName);
L
Liu Jicong 已提交
511

512 513
  int32_t len = 0;
  int32_t code = 0;
514
  tEncodeSize(tEncodeMqVgOffset, pOffset, len, code);
L
Liu Jicong 已提交
515
  if (code < 0) {
516
    return TSDB_CODE_INVALID_PARA;
L
Liu Jicong 已提交
517
  }
518

L
Liu Jicong 已提交
519
  void* buf = taosMemoryCalloc(1, sizeof(SMsgHead) + len);
L
Liu Jicong 已提交
520 521
  if (buf == NULL) {
    taosMemoryFree(pOffset);
522
    return TSDB_CODE_OUT_OF_MEMORY;
L
Liu Jicong 已提交
523
  }
524

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

L
Liu Jicong 已提交
527 528 529 530
  void* abuf = POINTER_SHIFT(buf, sizeof(SMsgHead));

  SEncoder encoder;
  tEncoderInit(&encoder, abuf, len);
531
  tEncodeMqVgOffset(&encoder, pOffset);
L
Liu Jicong 已提交
532
  tEncoderClear(&encoder);
L
Liu Jicong 已提交
533 534

  // build param
535
  SMqCommitCbParam* pParam = taosMemoryCalloc(1, sizeof(SMqCommitCbParam));
L
Liu Jicong 已提交
536
  if (pParam == NULL) {
L
Liu Jicong 已提交
537
    taosMemoryFree(pOffset);
L
Liu Jicong 已提交
538
    taosMemoryFree(buf);
539
    return TSDB_CODE_OUT_OF_MEMORY;
L
Liu Jicong 已提交
540
  }
541

L
Liu Jicong 已提交
542 543
  pParam->params = pParamSet;
  pParam->pOffset = pOffset;
H
Haojun Liao 已提交
544 545 546
  pParam->vgId = pVg->vgId;
  pParam->pTmq = tmq;

H
Haojun Liao 已提交
547
  tstrncpy(pParam->topicName, pTopicName, tListLen(pParam->topicName));
L
Liu Jicong 已提交
548 549 550 551

  // build send info
  SMsgSendInfo* pMsgSendInfo = taosMemoryCalloc(1, sizeof(SMsgSendInfo));
  if (pMsgSendInfo == NULL) {
L
Liu Jicong 已提交
552
    taosMemoryFree(pOffset);
L
Liu Jicong 已提交
553 554
    taosMemoryFree(buf);
    taosMemoryFree(pParam);
555
    return TSDB_CODE_OUT_OF_MEMORY;
L
Liu Jicong 已提交
556
  }
557

558
  pMsgSendInfo->msgInfo = (SDataBuf) { .pData = buf, .len = sizeof(SMsgHead) + len, .handle = NULL };
L
Liu Jicong 已提交
559 560 561 562

  pMsgSendInfo->requestId = generateRequestId();
  pMsgSendInfo->requestObjRefId = 0;
  pMsgSendInfo->param = pParam;
L
Liu Jicong 已提交
563
  pMsgSendInfo->paramFreeFp = taosMemoryFree;
564
  pMsgSendInfo->fp = tmqCommitCb;
565
  pMsgSendInfo->msgType = type;
L
Liu Jicong 已提交
566

L
Liu Jicong 已提交
567 568 569
  atomic_add_fetch_32(&pParamSet->waitingRspNum, 1);
  atomic_add_fetch_32(&pParamSet->totalRspNum, 1);

H
Haojun Liao 已提交
570
  SEp* pEp = GET_ACTIVE_EP(&pVg->epSet);
571
  char offsetBuf[TSDB_OFFSET_LEN] = {0};
572
  tFormatOffset(offsetBuf, tListLen(offsetBuf), &pOffset->offset.val);
573

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

L
Liu Jicong 已提交
580 581
  int64_t transporterId = 0;
  asyncSendMsgToServer(tmq->pTscObj->pAppInfo->pTransporter, &pVg->epSet, &transporterId, pMsgSendInfo);
582 583

  return TSDB_CODE_SUCCESS;
L
Liu Jicong 已提交
584 585
}

H
Haojun Liao 已提交
586 587 588 589 590 591 592 593 594 595 596
static SMqClientTopic* getTopicByName(tmq_t* tmq, const char* pTopicName) {
  int32_t numOfTopics = taosArrayGetSize(tmq->clientTopics);
  for (int32_t i = 0; i < numOfTopics; ++i) {
    SMqClientTopic* pTopic = taosArrayGet(tmq->clientTopics, i);
    if (strcmp(pTopic->topicName, pTopicName) != 0) {
      continue;
    }

    return pTopic;
  }

H
Haojun Liao 已提交
597
  tscError("consumer:0x%" PRIx64 ", total:%d, failed to find topic:%s", tmq->consumerId, numOfTopics, pTopicName);
H
Haojun Liao 已提交
598 599 600
  return NULL;
}

601
static void asyncCommitOffset(tmq_t* tmq, const TAOS_RES* pRes, int32_t type, tmq_commit_cb* pCommitFp, void* userParam) {
602 603 604 605 606 607 608 609 610 611 612 613
  char*   pTopicName = NULL;
  int32_t vgId = 0;
  int32_t code = 0;

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

  if (TD_RES_TMQ(pRes)) {
    SMqRspObj* pRspObj = (SMqRspObj*)pRes;
    pTopicName = pRspObj->topic;
L
Liu Jicong 已提交
614
    vgId = pRspObj->vgId;
615 616 617
  } else if (TD_RES_TMQ_META(pRes)) {
    SMqMetaRspObj* pMetaRspObj = (SMqMetaRspObj*)pRes;
    pTopicName = pMetaRspObj->topic;
L
Liu Jicong 已提交
618
    vgId = pMetaRspObj->vgId;
619 620 621
  } else if (TD_RES_TMQ_METADATA(pRes)) {
    SMqTaosxRspObj* pRspObj = (SMqTaosxRspObj*)pRes;
    pTopicName = pRspObj->topic;
622
    vgId = pRspObj->vgId;
L
Liu Jicong 已提交
623
  } else {
624 625
    pCommitFp(tmq, TSDB_CODE_TMQ_INVALID_MSG, userParam);
    return;
L
Liu Jicong 已提交
626 627 628 629
  }

  SMqCommitCbParamSet* pParamSet = taosMemoryCalloc(1, sizeof(SMqCommitCbParamSet));
  if (pParamSet == NULL) {
630 631
    pCommitFp(tmq, TSDB_CODE_OUT_OF_MEMORY, userParam);
    return;
L
Liu Jicong 已提交
632
  }
H
Haojun Liao 已提交
633

634 635
  pParamSet->refId = tmq->refId;
  pParamSet->epoch = tmq->epoch;
636
  pParamSet->callbackFn = pCommitFp;
L
Liu Jicong 已提交
637
  pParamSet->userParam = userParam;
L
Liu Jicong 已提交
638

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

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

H
Haojun Liao 已提交
643 644
  SMqClientTopic* pTopic = getTopicByName(tmq, pTopicName);
  if (pTopic == NULL) {
X
Xiaoyu Wang 已提交
645 646
    tscWarn("consumer:0x%" PRIx64 " failed to find the specified topic:%s, total topics:%d", tmq->consumerId,
            pTopicName, numOfTopics);
647 648 649 650
    taosMemoryFree(pParamSet);
    pCommitFp(tmq, TSDB_CODE_SUCCESS, userParam);
    return;
  }
L
Liu Jicong 已提交
651

652 653 654
  int32_t j = 0;
  int32_t numOfVgroups = taosArrayGetSize(pTopic->vgs);
  for (j = 0; j < numOfVgroups; j++) {
wmmhello's avatar
wmmhello 已提交
655
    SMqClientVg* pVg = (SMqClientVg*)taosArrayGet(pTopic->vgs, j);
656 657
    if (pVg->vgId == vgId) {
      break;
L
Liu Jicong 已提交
658
    }
L
Liu Jicong 已提交
659
  }
L
Liu Jicong 已提交
660

661
  if (j == numOfVgroups) {
X
Xiaoyu Wang 已提交
662 663
    tscWarn("consumer:0x%" PRIx64 " failed to find the specified vgId:%d, total Vgs:%d, topic:%s", tmq->consumerId,
            vgId, numOfVgroups, pTopicName);
L
Liu Jicong 已提交
664
    taosMemoryFree(pParamSet);
665 666
    pCommitFp(tmq, TSDB_CODE_SUCCESS, userParam);
    return;
L
Liu Jicong 已提交
667 668
  }

wmmhello's avatar
wmmhello 已提交
669
  SMqClientVg* pVg = (SMqClientVg*)taosArrayGet(pTopic->vgs, j);
H
Haojun Liao 已提交
670
  if (pVg->offsetInfo.currentOffset.type > 0 && !tOffsetEqual(&pVg->offsetInfo.currentOffset, &pVg->offsetInfo.committedOffset)) {
671
    code = doSendCommitMsg(tmq, pVg, pTopic->topicName, pParamSet, j, numOfVgroups, type);
L
Liu Jicong 已提交
672

673 674 675 676 677
    // failed to commit, callback user function directly.
    if (code != TSDB_CODE_SUCCESS) {
      taosMemoryFree(pParamSet);
      pCommitFp(tmq, code, userParam);
    }
X
Xiaoyu Wang 已提交
678
  } else {  // do not perform commit, callback user function directly.
679 680
    taosMemoryFree(pParamSet);
    pCommitFp(tmq, code, userParam);
L
Liu Jicong 已提交
681 682 683
  }
}

684
static void asyncCommitAllOffsets(tmq_t* tmq, tmq_commit_cb* pCommitFp, void* userParam) {
685 686
  SMqCommitCbParamSet* pParamSet = taosMemoryCalloc(1, sizeof(SMqCommitCbParamSet));
  if (pParamSet == NULL) {
687 688
    pCommitFp(tmq, TSDB_CODE_OUT_OF_MEMORY, userParam);
    return;
689
  }
690 691 692

  pParamSet->refId = tmq->refId;
  pParamSet->epoch = tmq->epoch;
693
  pParamSet->callbackFn = pCommitFp;
694 695
  pParamSet->userParam = userParam;

696 697 698
  // init as 1 to prevent concurrency issue
  pParamSet->waitingRspNum = 1;

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

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

706 707
    tscDebug("consumer:0x%" PRIx64 " commit offset for topics:%s, numOfVgs:%d", tmq->consumerId, pTopic->topicName,
             numOfVgroups);
708
    for (int32_t j = 0; j < numOfVgroups; j++) {
709 710
      SMqClientVg* pVg = taosArrayGet(pTopic->vgs, j);

H
Haojun Liao 已提交
711
      if (pVg->offsetInfo.currentOffset.type > 0 && !tOffsetEqual(&pVg->offsetInfo.currentOffset, &pVg->offsetInfo.committedOffset)) {
712
        int32_t code = doSendCommitMsg(tmq, pVg, pTopic->topicName, pParamSet, j, numOfVgroups, TDMT_VND_TMQ_COMMIT_OFFSET);
713 714
        if (code != TSDB_CODE_SUCCESS) {
          tscError("consumer:0x%" PRIx64 " topic:%s vgId:%d offset:%" PRId64 " failed, code:%s ordinal:%d/%d",
H
Haojun Liao 已提交
715
                   tmq->consumerId, pTopic->topicName, pVg->vgId, pVg->offsetInfo.committedOffset.version, tstrerror(terrno),
716
                   j + 1, numOfVgroups);
L
Liu Jicong 已提交
717 718
          continue;
        }
H
Haojun Liao 已提交
719 720

        // update the offset value.
H
Haojun Liao 已提交
721
        pVg->offsetInfo.committedOffset = pVg->offsetInfo.currentOffset;
722
      } else {
D
dapan1121 已提交
723
        tscDebug("consumer:0x%" PRIx64 " topic:%s vgId:%d, no commit, current:%" PRId64 ", ordinal:%d/%d",
H
Haojun Liao 已提交
724
                 tmq->consumerId, pTopic->topicName, pVg->vgId, pVg->offsetInfo.currentOffset.version, j + 1, numOfVgroups);
725 726 727 728
      }
    }
  }

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

L
Liu Jicong 已提交
732
  // no request is sent
L
Liu Jicong 已提交
733 734
  if (pParamSet->totalRspNum == 0) {
    taosMemoryFree(pParamSet);
735 736
    pCommitFp(tmq, TSDB_CODE_SUCCESS, userParam);
    return;
L
Liu Jicong 已提交
737 738
  }

L
Liu Jicong 已提交
739
  // count down since waiting rsp num init as 1
740
  commitRspCountDown(pParamSet, tmq->consumerId, "", 0);
741 742
}

743 744
static void generateTimedTask(int64_t refId, int32_t type) {
  tmq_t* tmq = taosAcquireRef(tmqMgmt.rsetId, refId);
wmmhello's avatar
wmmhello 已提交
745 746 747 748 749 750 751 752 753
  if(tmq == NULL) return;

  int8_t* pTaskType = taosAllocateQitem(sizeof(int8_t), DEF_QITEM, 0);
  if(pTaskType == NULL) return;

  *pTaskType = type;
  taosWriteQitem(tmq->delayedTask, pTaskType);
  tsem_post(&tmq->rspSem);
  taosReleaseRef(tmqMgmt.rsetId, refId);
754 755 756 757 758
}

void tmqAssignAskEpTask(void* param, void* tmrId) {
  int64_t refId = *(int64_t*)param;
  generateTimedTask(refId, TMQ_DELAYED_TASK__ASK_EP);
759
  taosMemoryFree(param);
L
Liu Jicong 已提交
760 761 762
}

void tmqAssignDelayedCommitTask(void* param, void* tmrId) {
763
  int64_t refId = *(int64_t*)param;
764
  generateTimedTask(refId, TMQ_DELAYED_TASK__COMMIT);
765
  taosMemoryFree(param);
L
Liu Jicong 已提交
766 767
}

wmmhello's avatar
wmmhello 已提交
768 769 770 771 772 773 774 775 776 777 778 779 780
//void tmqAssignDelayedReportTask(void* param, void* tmrId) {
//  int64_t refId = *(int64_t*)param;
//  tmq_t*  tmq = taosAcquireRef(tmqMgmt.rsetId, refId);
//  if (tmq != NULL) {
//    int8_t* pTaskType = taosAllocateQitem(sizeof(int8_t), DEF_QITEM, 0);
//    *pTaskType = TMQ_DELAYED_TASK__REPORT;
//    taosWriteQitem(tmq->delayedTask, pTaskType);
//    tsem_post(&tmq->rspSem);
//  }
//
//  taosReleaseRef(tmqMgmt.rsetId, refId);
//  taosMemoryFree(param);
//}
L
Liu Jicong 已提交
781

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

void tmqSendHbReq(void* param, void* tmrId) {
791
  int64_t refId = *(int64_t*)param;
792

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

  SMqHbReq req = {0};
  req.consumerId = tmq->consumerId;
  req.epoch = tmq->epoch;
wmmhello's avatar
wmmhello 已提交
802
//  if(tmq->needReportOffsetRows){
803 804 805 806 807 808 809 810 811 812 813 814
    req.topics = taosArrayInit(taosArrayGetSize(tmq->clientTopics), sizeof(TopicOffsetRows));
    for(int i = 0; i < taosArrayGetSize(tmq->clientTopics); i++){
      SMqClientTopic* pTopic = taosArrayGet(tmq->clientTopics, i);
      int32_t         numOfVgroups = taosArrayGetSize(pTopic->vgs);
      TopicOffsetRows* data = taosArrayReserve(req.topics, 1);
      strcpy(data->topicName, pTopic->topicName);
      data->offsetRows = taosArrayInit(numOfVgroups, sizeof(OffsetRows));
      for(int j = 0; j < numOfVgroups; j++){
        SMqClientVg* pVg = taosArrayGet(pTopic->vgs, j);
        OffsetRows* offRows = taosArrayReserve(data->offsetRows, 1);
        offRows->vgId = pVg->vgId;
        offRows->rows = pVg->numOfRows;
wmmhello's avatar
wmmhello 已提交
815
        offRows->offset = pVg->offsetInfo.currentOffset;
816 817
        char buf[TSDB_OFFSET_LEN] = {0};
        tFormatOffset(buf, TSDB_OFFSET_LEN, &offRows->offset);
wmmhello's avatar
wmmhello 已提交
818
        tscInfo("consumer:0x%" PRIx64 ",report offset: vgId:%d, offset:%s, rows:%"PRId64, tmq->consumerId, offRows->vgId, buf, offRows->rows);
819
      }
820
    }
wmmhello's avatar
wmmhello 已提交
821 822
//    tmq->needReportOffsetRows = false;
//  }
D
dapan1121 已提交
823

L
Liu Jicong 已提交
824
  int32_t tlen = tSerializeSMqHbReq(NULL, 0, &req);
D
dapan1121 已提交
825 826
  if (tlen < 0) {
    tscError("tSerializeSMqHbReq failed");
827
    goto OVER;
D
dapan1121 已提交
828
  }
829

L
Liu Jicong 已提交
830
  void* pReq = taosMemoryCalloc(1, tlen);
D
dapan1121 已提交
831 832
  if (tlen < 0) {
    tscError("failed to malloc MqHbReq msg, size:%d", tlen);
833
    goto OVER;
D
dapan1121 已提交
834
  }
835

D
dapan1121 已提交
836 837 838
  if (tSerializeSMqHbReq(pReq, tlen, &req) < 0) {
    tscError("tSerializeSMqHbReq %d failed", tlen);
    taosMemoryFree(pReq);
839
    goto OVER;
D
dapan1121 已提交
840
  }
841 842 843 844

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

848
  sendInfo->msgInfo = (SDataBuf){ .pData = pReq, .len = tlen, .handle = NULL };
849 850 851 852 853

  sendInfo->requestId = generateRequestId();
  sendInfo->requestObjRefId = 0;
  sendInfo->param = NULL;
  sendInfo->fp = tmqHbCb;
L
Liu Jicong 已提交
854
  sendInfo->msgType = TDMT_MND_TMQ_HB;
855 856 857 858 859 860 861

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

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

OVER:
862
  tDeatroySMqHbReq(&req);
863
  taosTmrReset(tmqSendHbReq, 1000, param, tmqMgmt.timer, &tmq->hbLiveTimer);
864
  taosReleaseRef(tmqMgmt.rsetId, refId);
865 866
}

867 868
static void defaultCommitCbFn(tmq_t* pTmq, int32_t code, void* param) {
  if (code != 0) {
869
    tscError("consumer:0x%" PRIx64 ", failed to commit offset, code:%s", pTmq->consumerId, tstrerror(code));
870 871 872
  }
}

873
int32_t tmqHandleAllDelayedTask(tmq_t* pTmq) {
L
Liu Jicong 已提交
874
  STaosQall* qall = taosAllocateQall();
875
  taosReadAllQitems(pTmq->delayedTask, qall);
L
Liu Jicong 已提交
876

877 878 879 880
  if (qall->numOfItems == 0) {
    taosFreeQall(qall);
    return TSDB_CODE_SUCCESS;
  }
881

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

886
  while (pTaskType != NULL) {
887
    if (*pTaskType == TMQ_DELAYED_TASK__ASK_EP) {
888
      asyncAskEp(pTmq, addToQueueCallbackFn, NULL);
889 890

      int64_t* pRefId = taosMemoryMalloc(sizeof(int64_t));
891
      *pRefId = pTmq->refId;
892

X
Xiaoyu Wang 已提交
893
      tscDebug("consumer:0x%" PRIx64 " retrieve ep from mnode in 1s", pTmq->consumerId);
894
      taosTmrReset(tmqAssignAskEpTask, 1000, pRefId, tmqMgmt.timer, &pTmq->epTimer);
L
Liu Jicong 已提交
895
    } else if (*pTaskType == TMQ_DELAYED_TASK__COMMIT) {
X
Xiaoyu Wang 已提交
896
      tmq_commit_cb* pCallbackFn = pTmq->commitCb ? pTmq->commitCb : defaultCommitCbFn;
897 898

      asyncCommitAllOffsets(pTmq, pCallbackFn, pTmq->commitCbUserParam);
899
      int64_t* pRefId = taosMemoryMalloc(sizeof(int64_t));
900
      *pRefId = pTmq->refId;
901

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

L
Liu Jicong 已提交
908
    taosFreeQitem(pTaskType);
909
    taosGetQitem(qall, (void**)&pTaskType);
L
Liu Jicong 已提交
910
  }
911

L
Liu Jicong 已提交
912 913 914 915
  taosFreeQall(qall);
  return 0;
}

916
static void* tmqFreeRspWrapper(SMqRspWrapper* rspWrapper) {
L
Liu Jicong 已提交
917 918 919 920 921 922 923
  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;
924 925
    taosMemoryFreeClear(pRsp->pEpset);

L
Liu Jicong 已提交
926 927 928
    taosArrayDestroyP(pRsp->dataRsp.blockData, taosMemoryFree);
    taosArrayDestroy(pRsp->dataRsp.blockDataLen);
    taosArrayDestroyP(pRsp->dataRsp.blockTbName, taosMemoryFree);
929
    taosArrayDestroyP(pRsp->dataRsp.blockSchema, (FDelete)tDeleteSchemaWrapper);
L
Liu Jicong 已提交
930 931
  } else if (rspWrapper->tmqRspType == TMQ_MSG_TYPE__POLL_META_RSP) {
    SMqPollRspWrapper* pRsp = (SMqPollRspWrapper*)rspWrapper;
932 933
    taosMemoryFreeClear(pRsp->pEpset);

L
Liu Jicong 已提交
934 935 936
    taosMemoryFree(pRsp->metaRsp.metaRsp);
  } else if (rspWrapper->tmqRspType == TMQ_MSG_TYPE__TAOSX_RSP) {
    SMqPollRspWrapper* pRsp = (SMqPollRspWrapper*)rspWrapper;
937 938
    taosMemoryFreeClear(pRsp->pEpset);

L
Liu Jicong 已提交
939 940 941
    taosArrayDestroyP(pRsp->taosxRsp.blockData, taosMemoryFree);
    taosArrayDestroy(pRsp->taosxRsp.blockDataLen);
    taosArrayDestroyP(pRsp->taosxRsp.blockTbName, taosMemoryFree);
942
    taosArrayDestroyP(pRsp->taosxRsp.blockSchema, (FDelete)tDeleteSchemaWrapper);
L
Liu Jicong 已提交
943 944 945 946
    // taosx
    taosArrayDestroy(pRsp->taosxRsp.createTableLen);
    taosArrayDestroyP(pRsp->taosxRsp.createTableReq, taosMemoryFree);
  }
947 948

  return NULL;
L
Liu Jicong 已提交
949 950
}

L
Liu Jicong 已提交
951
void tmqClearUnhandleMsg(tmq_t* tmq) {
L
Liu Jicong 已提交
952
  SMqRspWrapper* rspWrapper = NULL;
L
Liu Jicong 已提交
953
  while (1) {
L
Liu Jicong 已提交
954 955 956 957 958
    taosGetQitem(tmq->qall, (void**)&rspWrapper);
    if (rspWrapper) {
      tmqFreeRspWrapper(rspWrapper);
      taosFreeQitem(rspWrapper);
    } else {
L
Liu Jicong 已提交
959
      break;
L
Liu Jicong 已提交
960
    }
L
Liu Jicong 已提交
961 962
  }

L
Liu Jicong 已提交
963
  rspWrapper = NULL;
L
Liu Jicong 已提交
964 965
  taosReadAllQitems(tmq->mqueue, tmq->qall);
  while (1) {
L
Liu Jicong 已提交
966 967 968 969 970
    taosGetQitem(tmq->qall, (void**)&rspWrapper);
    if (rspWrapper) {
      tmqFreeRspWrapper(rspWrapper);
      taosFreeQitem(rspWrapper);
    } else {
L
Liu Jicong 已提交
971
      break;
L
Liu Jicong 已提交
972
    }
L
Liu Jicong 已提交
973 974 975
  }
}

D
dapan1121 已提交
976
int32_t tmqSubscribeCb(void* param, SDataBuf* pMsg, int32_t code) {
L
Liu Jicong 已提交
977 978
  SMqSubscribeCbParam* pParam = (SMqSubscribeCbParam*)param;
  pParam->rspErr = code;
dengyihao's avatar
dengyihao 已提交
979 980

  taosMemoryFree(pMsg->pEpSet);
L
Liu Jicong 已提交
981 982 983
  tsem_post(&pParam->rspSem);
  return 0;
}
984

L
Liu Jicong 已提交
985
int32_t tmq_subscription(tmq_t* tmq, tmq_list_t** topics) {
X
Xiaoyu Wang 已提交
986 987 988 989
  if (*topics == NULL) {
    *topics = tmq_list_new();
  }
  for (int i = 0; i < taosArrayGetSize(tmq->clientTopics); i++) {
L
Liu Jicong 已提交
990
    SMqClientTopic* topic = taosArrayGet(tmq->clientTopics, i);
L
Liu Jicong 已提交
991
    tmq_list_append(*topics, strchr(topic->topicName, '.') + 1);
X
Xiaoyu Wang 已提交
992
  }
L
Liu Jicong 已提交
993
  return 0;
X
Xiaoyu Wang 已提交
994 995
}

L
Liu Jicong 已提交
996
int32_t tmq_unsubscribe(tmq_t* tmq) {
997 998 999 1000 1001 1002 1003 1004
  if (tmq->autoCommit) {
    int32_t rsp = tmq_commit_sync(tmq, NULL);
    if (rsp != 0) {
      return rsp;
    }
  }
  taosSsleep(2);  // sleep 2s for hb to send offset and rows to server

L
Liu Jicong 已提交
1005 1006
  int32_t     rsp;
  int32_t     retryCnt = 0;
L
Liu Jicong 已提交
1007
  tmq_list_t* lst = tmq_list_new();
L
Liu Jicong 已提交
1008 1009 1010 1011 1012 1013 1014 1015 1016 1017
  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 已提交
1018 1019
  tmq_list_destroy(lst);
  return rsp;
X
Xiaoyu Wang 已提交
1020 1021
}

1022 1023 1024 1025 1026 1027
static void freeClientVgImpl(void* param) {
  SMqClientTopic* pTopic = param;
  taosMemoryFreeClear(pTopic->schema.pSchema);
  taosArrayDestroy(pTopic->vgs);
}

1028
void tmqFreeImpl(void* handle) {
1029 1030
  tmq_t*  tmq = (tmq_t*)handle;
  int64_t id = tmq->consumerId;
L
Liu Jicong 已提交
1031

1032
  // TODO stop timer
L
Liu Jicong 已提交
1033 1034 1035 1036
  if (tmq->mqueue) {
    tmqClearUnhandleMsg(tmq);
    taosCloseQueue(tmq->mqueue);
  }
L
Liu Jicong 已提交
1037

H
Haojun Liao 已提交
1038 1039 1040 1041 1042
  if (tmq->delayedTask) {
    taosCloseQueue(tmq->delayedTask);
  }

  taosFreeQall(tmq->qall);
1043
  tsem_destroy(&tmq->rspSem);
L
Liu Jicong 已提交
1044

1045
  taosArrayDestroyEx(tmq->clientTopics, freeClientVgImpl);
1046 1047
  taos_close_internal(tmq->pTscObj);
  taosMemoryFree(tmq);
1048 1049

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

1052 1053 1054 1055 1056 1057 1058 1059 1060
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);
1061
  if (tmqMgmt.rsetId < 0) {
1062 1063 1064 1065
    tmqInitRes = terrno;
  }
}

L
Liu Jicong 已提交
1066
tmq_t* tmq_consumer_new(tmq_conf_t* conf, char* errstr, int32_t errstrLen) {
1067 1068 1069 1070
  taosThreadOnce(&tmqInit, tmqMgmtInit);
  if (tmqInitRes != 0) {
    terrno = tmqInitRes;
    return NULL;
L
Liu Jicong 已提交
1071 1072
  }

L
Liu Jicong 已提交
1073 1074
  tmq_t* pTmq = taosMemoryCalloc(1, sizeof(tmq_t));
  if (pTmq == NULL) {
L
Liu Jicong 已提交
1075
    terrno = TSDB_CODE_OUT_OF_MEMORY;
1076
    tscError("failed to create consumer, groupId:%s, code:%s", conf->groupId, terrstr());
L
Liu Jicong 已提交
1077 1078
    return NULL;
  }
L
Liu Jicong 已提交
1079

L
Liu Jicong 已提交
1080 1081 1082
  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 已提交
1083 1084 1085
  pTmq->clientTopics = taosArrayInit(0, sizeof(SMqClientTopic));
  pTmq->mqueue = taosOpenQueue();
  pTmq->delayedTask = taosOpenQueue();
H
Haojun Liao 已提交
1086
  pTmq->qall = taosAllocateQall();
L
Liu Jicong 已提交
1087

X
Xiaoyu Wang 已提交
1088 1089
  if (pTmq->clientTopics == NULL || pTmq->mqueue == NULL || pTmq->qall == NULL || pTmq->delayedTask == NULL ||
      conf->groupId[0] == 0) {
L
Liu Jicong 已提交
1090
    terrno = TSDB_CODE_OUT_OF_MEMORY;
X
Xiaoyu Wang 已提交
1091
    tscError("consumer:0x%" PRIx64 " setup failed since %s, groupId:%s", pTmq->consumerId, terrstr(), pTmq->groupId);
1092
    goto _failed;
L
Liu Jicong 已提交
1093
  }
L
Liu Jicong 已提交
1094

L
Liu Jicong 已提交
1095 1096
  // init status
  pTmq->status = TMQ_CONSUMER_STATUS__INIT;
L
Liu Jicong 已提交
1097 1098
  pTmq->pollCnt = 0;
  pTmq->epoch = 0;
wmmhello's avatar
wmmhello 已提交
1099
//  pTmq->needReportOffsetRows = true;
L
Liu Jicong 已提交
1100

L
Liu Jicong 已提交
1101 1102 1103
  // set conf
  strcpy(pTmq->clientId, conf->clientId);
  strcpy(pTmq->groupId, conf->groupId);
1104
  pTmq->withTbName = conf->withTbName;
L
Liu Jicong 已提交
1105
  pTmq->useSnapshot = conf->snapEnable;
L
Liu Jicong 已提交
1106
  pTmq->autoCommit = conf->autoCommit;
L
Liu Jicong 已提交
1107
  pTmq->autoCommitInterval = conf->autoCommitInterval;
L
Liu Jicong 已提交
1108 1109
  pTmq->commitCb = conf->commitCb;
  pTmq->commitCbUserParam = conf->commitCbUserParam;
L
Liu Jicong 已提交
1110
  pTmq->resetOffsetCfg = conf->resetOffset;
wmmhello's avatar
wmmhello 已提交
1111
  taosInitRWLatch(&pTmq->lock);
L
Liu Jicong 已提交
1112

1113 1114
  pTmq->hbBgEnable = conf->hbBgEnable;

L
Liu Jicong 已提交
1115
  // assign consumerId
L
Liu Jicong 已提交
1116
  pTmq->consumerId = tGenIdPI64();
X
Xiaoyu Wang 已提交
1117

L
Liu Jicong 已提交
1118 1119
  // init semaphore
  if (tsem_init(&pTmq->rspSem, 0, 0) != 0) {
1120
    tscError("consumer:0x %" PRIx64 " setup failed since %s, consumer group %s", pTmq->consumerId, terrstr(),
S
Shengliang Guan 已提交
1121
             pTmq->groupId);
1122
    goto _failed;
L
Liu Jicong 已提交
1123
  }
L
Liu Jicong 已提交
1124

L
Liu Jicong 已提交
1125 1126 1127
  // init connection
  pTmq->pTscObj = taos_connect_internal(conf->ip, user, pass, NULL, NULL, conf->port, CONN_TYPE__TMQ);
  if (pTmq->pTscObj == NULL) {
1128
    tscError("consumer:0x%" PRIx64 " setup failed since %s, groupId:%s", pTmq->consumerId, terrstr(), pTmq->groupId);
L
Liu Jicong 已提交
1129
    tsem_destroy(&pTmq->rspSem);
1130
    goto _failed;
L
Liu Jicong 已提交
1131
  }
L
Liu Jicong 已提交
1132

1133 1134
  pTmq->refId = taosAddRef(tmqMgmt.rsetId, pTmq);
  if (pTmq->refId < 0) {
1135
    goto _failed;
1136 1137
  }

1138
  if (pTmq->hbBgEnable) {
L
Liu Jicong 已提交
1139 1140
    int64_t* pRefId = taosMemoryMalloc(sizeof(int64_t));
    *pRefId = pTmq->refId;
1141
    pTmq->hbLiveTimer = taosTmrStart(tmqSendHbReq, 1000, pRefId, tmqMgmt.timer);
1142 1143
  }

1144
  char         buf[TSDB_OFFSET_LEN] = {0};
1145 1146
  STqOffsetVal offset = {.type = pTmq->resetOffsetCfg};
  tFormatOffset(buf, tListLen(buf), &offset);
X
Xiaoyu Wang 已提交
1147 1148 1149 1150
  tscInfo("consumer:0x%" PRIx64 " is setup, refId:%" PRId64
          ", groupId:%s, snapshot:%d, autoCommit:%d, commitInterval:%dms, offset:%s, backgroudHB:%d",
          pTmq->consumerId, pTmq->refId, pTmq->groupId, pTmq->useSnapshot, pTmq->autoCommit, pTmq->autoCommitInterval,
          buf, pTmq->hbBgEnable);
L
Liu Jicong 已提交
1151

1152
  return pTmq;
1153

1154 1155
_failed:
  tmqFreeImpl(pTmq);
L
Liu Jicong 已提交
1156
  return NULL;
1157 1158
}

L
Liu Jicong 已提交
1159
int32_t tmq_subscribe(tmq_t* tmq, const tmq_list_t* topic_list) {
1160
  const int32_t   MAX_RETRY_COUNT = 120 * 2;  // let's wait for 2 mins at most
L
Liu Jicong 已提交
1161 1162 1163
  const SArray*   container = &topic_list->container;
  int32_t         sz = taosArrayGetSize(container);
  void*           buf = NULL;
L
Liu Jicong 已提交
1164
  SMsgSendInfo*   sendInfo = NULL;
L
Liu Jicong 已提交
1165
  SCMSubscribeReq req = {0};
1166
  int32_t         code = 0;
1167

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

1170
  req.consumerId = tmq->consumerId;
L
Liu Jicong 已提交
1171
  tstrncpy(req.clientId, tmq->clientId, 256);
L
Liu Jicong 已提交
1172
  tstrncpy(req.cgroup, tmq->groupId, TSDB_CGROUP_LEN);
1173 1174
  req.topicNames = taosArrayInit(sz, sizeof(void*));

1175 1176 1177 1178
  if (req.topicNames == NULL) {
    code = TSDB_CODE_OUT_OF_MEMORY;
    goto FAIL;
  }
1179

1180 1181 1182 1183 1184
  req.withTbName = tmq->withTbName;
  req.autoCommit = tmq->autoCommit;
  req.autoCommitInterval = tmq->autoCommitInterval;
  req.resetOffsetCfg = tmq->resetOffsetCfg;

L
Liu Jicong 已提交
1185 1186
  for (int32_t i = 0; i < sz; i++) {
    char* topic = taosArrayGetP(container, i);
1187 1188

    SName name = {0};
L
Liu Jicong 已提交
1189 1190 1191 1192
    tNameSetDbName(&name, tmq->pTscObj->acctId, topic, strlen(topic));
    char* topicFName = taosMemoryCalloc(1, TSDB_TOPIC_FNAME_LEN);
    if (topicFName == NULL) {
      goto FAIL;
1193 1194
    }

1195
    tNameExtractFullName(&name, topicFName);
1196
    tscInfo("consumer:0x%" PRIx64 " subscribe topic:%s", tmq->consumerId, topicFName);
L
Liu Jicong 已提交
1197 1198

    taosArrayPush(req.topicNames, &topicFName);
1199 1200
  }

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

L
Liu Jicong 已提交
1203
  buf = taosMemoryMalloc(tlen);
1204 1205 1206 1207
  if (buf == NULL) {
    code = TSDB_CODE_OUT_OF_MEMORY;
    goto FAIL;
  }
L
Liu Jicong 已提交
1208

1209 1210 1211
  void* abuf = buf;
  tSerializeSCMSubscribeReq(&abuf, &req);

L
Liu Jicong 已提交
1212
  sendInfo = taosMemoryCalloc(1, sizeof(SMsgSendInfo));
1213 1214 1215 1216
  if (sendInfo == NULL) {
    code = TSDB_CODE_OUT_OF_MEMORY;
    goto FAIL;
  }
1217

H
Haojun Liao 已提交
1218
  SMqSubscribeCbParam param = { .rspErr = 0};
1219
  if (tsem_init(&param.rspSem, 0, 0) != 0) {
wmmhello's avatar
wmmhello 已提交
1220
    code = TSDB_CODE_TSC_INTERNAL_ERROR;
1221 1222
    goto FAIL;
  }
L
Liu Jicong 已提交
1223

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

L
Liu Jicong 已提交
1226 1227
  sendInfo->requestId = generateRequestId();
  sendInfo->requestObjRefId = 0;
L
Liu Jicong 已提交
1228 1229
  sendInfo->param = &param;
  sendInfo->fp = tmqSubscribeCb;
L
Liu Jicong 已提交
1230
  sendInfo->msgType = TDMT_MND_TMQ_SUBSCRIBE;
L
Liu Jicong 已提交
1231

1232 1233 1234 1235 1236
  SEpSet epSet = getEpSet_s(&tmq->pTscObj->pAppInfo->mgmtEp);

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

L
Liu Jicong 已提交
1237 1238
  // avoid double free if msg is sent
  buf = NULL;
L
Liu Jicong 已提交
1239
  sendInfo = NULL;
L
Liu Jicong 已提交
1240

L
Liu Jicong 已提交
1241 1242
  tsem_wait(&param.rspSem);
  tsem_destroy(&param.rspSem);
1243

1244 1245 1246 1247
  if (param.rspErr != 0) {
    code = param.rspErr;
    goto FAIL;
  }
L
Liu Jicong 已提交
1248

L
Liu Jicong 已提交
1249
  int32_t retryCnt = 0;
1250
  while (TSDB_CODE_MND_CONSUMER_NOT_READY == doAskEp(tmq)) {
1251
    if (retryCnt++ > MAX_RETRY_COUNT) {
wmmhello's avatar
wmmhello 已提交
1252
      tscError("consumer:0x%" PRIx64 ", mnd not ready for subscribe, retry:%d in 500ms", tmq->consumerId, retryCnt);
1253
      code = TSDB_CODE_MND_CONSUMER_NOT_READY;
L
Liu Jicong 已提交
1254 1255
      goto FAIL;
    }
1256

1257
    tscInfo("consumer:0x%" PRIx64 ", mnd not ready for subscribe, retry:%d in 500ms", tmq->consumerId, retryCnt);
L
Liu Jicong 已提交
1258 1259
    taosMsleep(500);
  }
1260

1261 1262
  // init ep timer
  if (tmq->epTimer == NULL) {
1263 1264 1265
    int64_t* pRefId1 = taosMemoryMalloc(sizeof(int64_t));
    *pRefId1 = tmq->refId;
    tmq->epTimer = taosTmrStart(tmqAssignAskEpTask, 1000, pRefId1, tmqMgmt.timer);
1266
  }
L
Liu Jicong 已提交
1267 1268

  // init auto commit timer
1269
  if (tmq->autoCommit && tmq->commitTimer == NULL) {
1270 1271 1272
    int64_t* pRefId2 = taosMemoryMalloc(sizeof(int64_t));
    *pRefId2 = tmq->refId;
    tmq->commitTimer = taosTmrStart(tmqAssignDelayedCommitTask, tmq->autoCommitInterval, pRefId2, tmqMgmt.timer);
L
Liu Jicong 已提交
1273 1274
  }

L
Liu Jicong 已提交
1275
FAIL:
L
Liu Jicong 已提交
1276
  taosArrayDestroyP(req.topicNames, taosMemoryFree);
L
Liu Jicong 已提交
1277
  taosMemoryFree(buf);
L
Liu Jicong 已提交
1278
  taosMemoryFree(sendInfo);
L
Liu Jicong 已提交
1279

L
Liu Jicong 已提交
1280
  return code;
1281 1282
}

L
Liu Jicong 已提交
1283
void tmq_conf_set_auto_commit_cb(tmq_conf_t* conf, tmq_commit_cb* cb, void* param) {
1284
  conf->commitCb = cb;
L
Liu Jicong 已提交
1285
  conf->commitCbUserParam = param;
L
Liu Jicong 已提交
1286
}
1287

wmmhello's avatar
wmmhello 已提交
1288 1289 1290 1291 1292 1293 1294 1295 1296 1297 1298 1299 1300 1301 1302 1303 1304 1305 1306 1307 1308 1309 1310 1311 1312 1313 1314 1315
static SMqClientVg* getVgInfo(tmq_t* tmq, char* topicName, int32_t  vgId){
  int32_t topicNumCur = taosArrayGetSize(tmq->clientTopics);
  for(int i = 0; i < topicNumCur; i++){
    SMqClientTopic* pTopicCur = taosArrayGet(tmq->clientTopics, i);
    if(strcmp(pTopicCur->topicName, topicName) == 0){
      int32_t vgNumCur = taosArrayGetSize(pTopicCur->vgs);
      for (int32_t j = 0; j < vgNumCur; j++) {
        SMqClientVg* pVgCur = taosArrayGet(pTopicCur->vgs, j);
        if(pVgCur->vgId == vgId){
          return pVgCur;
        }
      }
    }
  }
  return NULL;
}

static SMqClientTopic* getTopicInfo(tmq_t* tmq, char* topicName){
  int32_t topicNumCur = taosArrayGetSize(tmq->clientTopics);
  for(int i = 0; i < topicNumCur; i++){
    SMqClientTopic* pTopicCur = taosArrayGet(tmq->clientTopics, i);
    if(strcmp(pTopicCur->topicName, topicName) == 0){
      return pTopicCur;
    }
  }
  return NULL;
}

D
dapan1121 已提交
1316
int32_t tmqPollCb(void* param, SDataBuf* pMsg, int32_t code) {
X
Xiaoyu Wang 已提交
1317
  SMqPollCbParam* pParam = (SMqPollCbParam*)param;
1318 1319

  int64_t         refId = pParam->refId;
wmmhello's avatar
wmmhello 已提交
1320 1321
//  SMqClientVg*    pVg = pParam->pVg;
//  SMqClientTopic* pTopic = pParam->pTopic;
1322

1323
  tmq_t* tmq = taosAcquireRef(tmqMgmt.rsetId, refId);
1324 1325 1326
  if (tmq == NULL) {
    taosMemoryFree(pParam);
    taosMemoryFree(pMsg->pData);
dengyihao's avatar
dengyihao 已提交
1327
    taosMemoryFree(pMsg->pEpSet);
1328 1329 1330 1331
    terrno = TSDB_CODE_TMQ_CONSUMER_CLOSED;
    return -1;
  }

H
Haojun Liao 已提交
1332 1333 1334 1335
  int32_t  epoch = pParam->epoch;
  int32_t  vgId = pParam->vgId;
  uint64_t requestId = pParam->requestId;

L
Liu Jicong 已提交
1336
  if (code != 0) {
L
Liu Jicong 已提交
1337
    if (pMsg->pData) taosMemoryFree(pMsg->pData);
dengyihao's avatar
dengyihao 已提交
1338 1339
    if (pMsg->pEpSet) taosMemoryFree(pMsg->pEpSet);

H
Haojun Liao 已提交
1340
    // in case of consumer mismatch, wait for 500ms and retry
L
Liu Jicong 已提交
1341
    if (code == TSDB_CODE_TMQ_CONSUMER_MISMATCH) {
1342
//      taosMsleep(500);
L
Liu Jicong 已提交
1343
      atomic_store_8(&tmq->status, TMQ_CONSUMER_STATUS__RECOVER);
X
Xiaoyu Wang 已提交
1344 1345
      tscDebug("consumer:0x%" PRIx64 " wait for the re-balance, wait for 500ms and set status to be RECOVER",
               tmq->consumerId);
H
Haojun Liao 已提交
1346
    } else if (code == TSDB_CODE_TQ_NO_COMMITTED_OFFSET) {
S
Shengliang Guan 已提交
1347
      SMqPollRspWrapper* pRspWrapper = taosAllocateQitem(sizeof(SMqPollRspWrapper), DEF_QITEM, 0);
L
Liu Jicong 已提交
1348
      if (pRspWrapper == NULL) {
H
Haojun Liao 已提交
1349 1350
        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 已提交
1351 1352
        goto CREATE_MSG_FAIL;
      }
H
Haojun Liao 已提交
1353

L
Liu Jicong 已提交
1354 1355
      pRspWrapper->tmqRspType = TMQ_MSG_TYPE__END_RSP;
      taosWriteQitem(tmq->mqueue, pRspWrapper);
1356 1357
//    } else if (code == TSDB_CODE_WAL_LOG_NOT_EXIST) {  // poll data while insert
//      taosMsleep(5);
wmmhello's avatar
wmmhello 已提交
1358 1359 1360
    } else{
      tscError("consumer:0x%" PRIx64 " msg from vgId:%d discarded, epoch %d, since %s, reqId:0x%" PRIx64, tmq->consumerId,
               vgId, epoch, tstrerror(code), requestId);
L
Liu Jicong 已提交
1361
    }
H
Haojun Liao 已提交
1362

L
fix txn  
Liu Jicong 已提交
1363
    goto CREATE_MSG_FAIL;
L
Liu Jicong 已提交
1364 1365
  }

X
Xiaoyu Wang 已提交
1366
  int32_t msgEpoch = ((SMqRspHead*)pMsg->pData)->epoch;
1367 1368
  int32_t clientEpoch = atomic_load_32(&tmq->epoch);
  if (msgEpoch < clientEpoch) {
L
Liu Jicong 已提交
1369
    // do not write into queue since updating epoch reset
X
Xiaoyu Wang 已提交
1370 1371
    tscWarn("consumer:0x%" PRIx64
            " msg discard from vgId:%d since from earlier epoch, rsp epoch %d, current epoch %d, reqId:0x%" PRIx64,
1372
            tmq->consumerId, vgId, msgEpoch, clientEpoch, requestId);
H
Haojun Liao 已提交
1373

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

L
Liu Jicong 已提交
1377
    taosMemoryFree(pMsg->pData);
dengyihao's avatar
dengyihao 已提交
1378
    taosMemoryFree(pMsg->pEpSet);
wmmhello's avatar
wmmhello 已提交
1379 1380
    taosMemoryFree(pParam);

X
Xiaoyu Wang 已提交
1381 1382 1383
    return 0;
  }

1384
  if (msgEpoch != clientEpoch) {
H
Haojun Liao 已提交
1385
    tscWarn("consumer:0x%" PRIx64 " mismatch rsp from vgId:%d, epoch %d, current epoch %d, reqId:0x%" PRIx64,
1386
            tmq->consumerId, vgId, msgEpoch, clientEpoch, requestId);
X
Xiaoyu Wang 已提交
1387 1388
  }

L
Liu Jicong 已提交
1389 1390 1391
  // handle meta rsp
  int8_t rspType = ((SMqRspHead*)pMsg->pData)->mqMsgType;

S
Shengliang Guan 已提交
1392
  SMqPollRspWrapper* pRspWrapper = taosAllocateQitem(sizeof(SMqPollRspWrapper), DEF_QITEM, 0);
L
Liu Jicong 已提交
1393
  if (pRspWrapper == NULL) {
L
Liu Jicong 已提交
1394
    taosMemoryFree(pMsg->pData);
dengyihao's avatar
dengyihao 已提交
1395
    taosMemoryFree(pMsg->pEpSet);
X
Xiaoyu Wang 已提交
1396 1397
    tscWarn("consumer:0x%" PRIx64 " msg discard from vgId:%d, epoch %d since out of memory", tmq->consumerId, vgId,
            epoch);
L
fix txn  
Liu Jicong 已提交
1398
    goto CREATE_MSG_FAIL;
L
Liu Jicong 已提交
1399
  }
L
Liu Jicong 已提交
1400

L
Liu Jicong 已提交
1401
  pRspWrapper->tmqRspType = rspType;
wmmhello's avatar
wmmhello 已提交
1402 1403
//  pRspWrapper->vgHandle = pVg;
//  pRspWrapper->topicHandle = pTopic;
H
Haojun Liao 已提交
1404
  pRspWrapper->reqId = requestId;
1405
  pRspWrapper->pEpset = pMsg->pEpSet;
wmmhello's avatar
wmmhello 已提交
1406 1407
  pRspWrapper->vgId = vgId;
  strcpy(pRspWrapper->topicName, pParam->topicName);
L
Liu Jicong 已提交
1408

1409
  pMsg->pEpSet = NULL;
L
Liu Jicong 已提交
1410
  if (rspType == TMQ_MSG_TYPE__POLL_RSP) {
L
Liu Jicong 已提交
1411 1412
    SDecoder decoder;
    tDecoderInit(&decoder, POINTER_SHIFT(pMsg->pData, sizeof(SMqRspHead)), pMsg->len - sizeof(SMqRspHead));
1413
    tDecodeMqDataRsp(&decoder, &pRspWrapper->dataRsp);
wmmhello's avatar
wmmhello 已提交
1414
    tDecoderClear(&decoder);
L
Liu Jicong 已提交
1415
    memcpy(&pRspWrapper->dataRsp, pMsg->pData, sizeof(SMqRspHead));
L
Liu Jicong 已提交
1416

1417 1418
    char buf[TSDB_OFFSET_LEN];
    tFormatOffset(buf, TSDB_OFFSET_LEN, &pRspWrapper->dataRsp.rspOffset);
H
Haojun Liao 已提交
1419
    tscDebug("consumer:0x%" PRIx64 " recv poll rsp, vgId:%d, req ver:%" PRId64 ", rsp:%s type %d, reqId:0x%" PRIx64,
H
Haojun Liao 已提交
1420
             tmq->consumerId, vgId, pRspWrapper->dataRsp.reqOffset.version, buf, rspType, requestId);
L
Liu Jicong 已提交
1421
  } else if (rspType == TMQ_MSG_TYPE__POLL_META_RSP) {
1422 1423
    SDecoder decoder;
    tDecoderInit(&decoder, POINTER_SHIFT(pMsg->pData, sizeof(SMqRspHead)), pMsg->len - sizeof(SMqRspHead));
1424
    tDecodeMqMetaRsp(&decoder, &pRspWrapper->metaRsp);
1425
    tDecoderClear(&decoder);
L
Liu Jicong 已提交
1426
    memcpy(&pRspWrapper->metaRsp, pMsg->pData, sizeof(SMqRspHead));
L
Liu Jicong 已提交
1427 1428 1429 1430 1431 1432
  } 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));
X
Xiaoyu Wang 已提交
1433 1434
  } else {  // invalid rspType
    tscError("consumer:0x%" PRIx64 " invalid rsp msg received, type:%d ignored", tmq->consumerId, rspType);
L
Liu Jicong 已提交
1435
  }
L
Liu Jicong 已提交
1436

L
Liu Jicong 已提交
1437
  taosMemoryFree(pMsg->pData);
H
Haojun Liao 已提交
1438
  taosWriteQitem(tmq->mqueue, pRspWrapper);
L
Liu Jicong 已提交
1439

1440
  int32_t total = taosQueueItemSize(tmq->mqueue);
H
Haojun Liao 已提交
1441
  tscDebug("consumer:0x%" PRIx64 " put poll res into mqueue, type:%d, vgId:%d, total in queue:%d, reqId:0x%" PRIx64,
1442
           tmq->consumerId, rspType, vgId, total, requestId);
H
Haojun Liao 已提交
1443

1444
  tsem_post(&tmq->rspSem);
1445
  taosReleaseRef(tmqMgmt.rsetId, refId);
wmmhello's avatar
wmmhello 已提交
1446
  taosMemoryFree(pParam);
1447

L
Liu Jicong 已提交
1448
  return 0;
H
Haojun Liao 已提交
1449

L
fix txn  
Liu Jicong 已提交
1450
CREATE_MSG_FAIL:
L
Liu Jicong 已提交
1451
  if (epoch == tmq->epoch) {
wmmhello's avatar
wmmhello 已提交
1452 1453 1454 1455 1456
    taosWLockLatch(&tmq->lock);
    SMqClientVg* pVg = getVgInfo(tmq, pParam->topicName, vgId);
    if(pVg){
      atomic_store_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE);
    }
wmmhello's avatar
wmmhello 已提交
1457
    taosWUnLockLatch(&tmq->lock);
L
Liu Jicong 已提交
1458
  }
H
Haojun Liao 已提交
1459

1460
  tsem_post(&tmq->rspSem);
1461
  taosReleaseRef(tmqMgmt.rsetId, refId);
wmmhello's avatar
wmmhello 已提交
1462
  taosMemoryFree(pParam);
1463

L
Liu Jicong 已提交
1464
  return -1;
1465 1466
}

H
Haojun Liao 已提交
1467
typedef struct SVgroupSaveInfo {
wmmhello's avatar
wmmhello 已提交
1468 1469
  STqOffsetVal currentOffset;
  STqOffsetVal commitOffset;
H
Haojun Liao 已提交
1470 1471 1472
  int64_t      numOfRows;
} SVgroupSaveInfo;

H
Haojun Liao 已提交
1473 1474 1475 1476 1477 1478
static void initClientTopicFromRsp(SMqClientTopic* pTopic, SMqSubTopicEp* pTopicEp, SHashObj* pVgOffsetHashMap,
                                   tmq_t* tmq) {
  pTopic->schema = pTopicEp->schema;
  pTopicEp->schema.nCols = 0;
  pTopicEp->schema.pSchema = NULL;

X
Xiaoyu Wang 已提交
1479
  char    vgKey[TSDB_TOPIC_FNAME_LEN + 22];
H
Haojun Liao 已提交
1480 1481 1482 1483 1484
  int32_t vgNumGet = taosArrayGetSize(pTopicEp->vgs);

  tstrncpy(pTopic->topicName, pTopicEp->topic, TSDB_TOPIC_FNAME_LEN);
  tstrncpy(pTopic->db, pTopicEp->db, TSDB_DB_FNAME_LEN);

1485
  tscInfo("consumer:0x%" PRIx64 ", update topic:%s, new numOfVgs:%d", tmq->consumerId, pTopic->topicName, vgNumGet);
H
Haojun Liao 已提交
1486 1487 1488 1489
  pTopic->vgs = taosArrayInit(vgNumGet, sizeof(SMqClientVg));

  for (int32_t j = 0; j < vgNumGet; j++) {
    SMqSubVgEp* pVgEp = taosArrayGet(pTopicEp->vgs, j);
H
Haojun Liao 已提交
1490 1491

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

wmmhello's avatar
wmmhello 已提交
1494 1495
    STqOffsetVal offsetNew = {0};
    offsetNew.type = tmq->resetOffsetCfg;
H
Haojun Liao 已提交
1496 1497 1498 1499 1500

    SMqClientVg clientVg = {
        .pollCnt = 0,
        .vgId = pVgEp->vgId,
        .epSet = pVgEp->epSet,
wmmhello's avatar
wmmhello 已提交
1501
        .vgStatus = TMQ_VG_STATUS__IDLE,
H
Haojun Liao 已提交
1502
        .vgSkipCnt = 0,
H
Haojun Liao 已提交
1503
        .emptyBlockReceiveTs = 0,
wmmhello's avatar
wmmhello 已提交
1504
        .numOfRows = pInfo ? pInfo->numOfRows : 0,
H
Haojun Liao 已提交
1505 1506
    };

wmmhello's avatar
wmmhello 已提交
1507 1508
    clientVg.offsetInfo.currentOffset = pInfo ? pInfo->currentOffset : offsetNew;
    clientVg.offsetInfo.committedOffset = pInfo ? pInfo->commitOffset : offsetNew;
H
Haojun Liao 已提交
1509 1510
    clientVg.offsetInfo.walVerBegin = -1;
    clientVg.offsetInfo.walVerEnd = -1;
1511 1512 1513
    clientVg.seekUpdated = false;
    clientVg.receivedInfoFromVnode = false;

H
Haojun Liao 已提交
1514 1515 1516 1517 1518 1519 1520 1521 1522 1523 1524 1525 1526
    taosArrayPush(pTopic->vgs, &clientVg);
  }
}

static void freeClientVgInfo(void* param) {
  SMqClientTopic* pTopic = param;
  if (pTopic->schema.nCols) {
    taosMemoryFreeClear(pTopic->schema.pSchema);
  }

  taosArrayDestroy(pTopic->vgs);
}

1527
static bool doUpdateLocalEp(tmq_t* tmq, int32_t epoch, const SMqAskEpRsp* pRsp) {
1528 1529
  bool set = false;

1530
  int32_t topicNumCur = taosArrayGetSize(tmq->clientTopics);
1531
  int32_t topicNumGet = taosArrayGetSize(pRsp->topics);
1532

X
Xiaoyu Wang 已提交
1533
  char vgKey[TSDB_TOPIC_FNAME_LEN + 22];
1534
  tscInfo("consumer:0x%" PRIx64 " update ep epoch from %d to epoch %d, incoming topics:%d, existed topics:%d",
1535
           tmq->consumerId, tmq->epoch, epoch, topicNumGet, topicNumCur);
wmmhello's avatar
wmmhello 已提交
1536 1537 1538
  if (epoch <= tmq->epoch) {
    return false;
  }
1539 1540 1541 1542 1543 1544

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

H
Haojun Liao 已提交
1545 1546
  SHashObj* pVgOffsetHashMap = taosHashInit(64, MurmurHash3_32, false, HASH_NO_LOCK);
  if (pVgOffsetHashMap == NULL) {
1547 1548 1549
    taosArrayDestroy(newTopics);
    return false;
  }
1550

H
Haojun Liao 已提交
1551
  // todo extract method
1552 1553 1554 1555 1556
  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);
1557
      tscInfo("consumer:0x%" PRIx64 ", current vg num: %d", tmq->consumerId, vgNumCur);
1558 1559
      for (int32_t j = 0; j < vgNumCur; j++) {
        SMqClientVg* pVgCur = taosArrayGet(pTopicCur->vgs, j);
H
Haojun Liao 已提交
1560 1561
        makeTopicVgroupKey(vgKey, pTopicCur->topicName, pVgCur->vgId);

1562 1563
        char buf[TSDB_OFFSET_LEN];
        tFormatOffset(buf, TSDB_OFFSET_LEN, &pVgCur->offsetInfo.currentOffset);
1564
        tscInfo("consumer:0x%" PRIx64 ", epoch:%d vgId:%d vgKey:%s, offset:%s", tmq->consumerId, epoch, pVgCur->vgId,
X
Xiaoyu Wang 已提交
1565
                 vgKey, buf);
H
Haojun Liao 已提交
1566

wmmhello's avatar
wmmhello 已提交
1567
        SVgroupSaveInfo info = {.currentOffset = pVgCur->offsetInfo.currentOffset, .commitOffset = pVgCur->offsetInfo.committedOffset, .numOfRows = pVgCur->numOfRows};
H
Haojun Liao 已提交
1568
        taosHashPut(pVgOffsetHashMap, vgKey, strlen(vgKey), &info, sizeof(SVgroupSaveInfo));
1569 1570 1571 1572 1573 1574 1575
      }
    }
  }

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

H
Haojun Liao 已提交
1580 1581
  taosHashCleanup(pVgOffsetHashMap);

wmmhello's avatar
wmmhello 已提交
1582
  taosWLockLatch(&tmq->lock);
1583
  // destroy current buffered existed topics info
1584
  if (tmq->clientTopics) {
H
Haojun Liao 已提交
1585
    taosArrayDestroyEx(tmq->clientTopics, freeClientVgInfo);
X
Xiaoyu Wang 已提交
1586
  }
H
Haojun Liao 已提交
1587
  tmq->clientTopics = newTopics;
wmmhello's avatar
wmmhello 已提交
1588
  taosWUnLockLatch(&tmq->lock);
1589

X
Xiaoyu Wang 已提交
1590
  int8_t flag = (topicNumGet == 0) ? TMQ_CONSUMER_STATUS__NO_TOPIC : TMQ_CONSUMER_STATUS__READY;
H
Haojun Liao 已提交
1591
  atomic_store_8(&tmq->status, flag);
X
Xiaoyu Wang 已提交
1592
  atomic_store_32(&tmq->epoch, epoch);
H
Haojun Liao 已提交
1593

1594
  tscInfo("consumer:0x%" PRIx64 " update topic info completed", tmq->consumerId);
X
Xiaoyu Wang 已提交
1595 1596 1597
  return set;
}

1598
int32_t askEpCallbackFn(void* param, SDataBuf* pMsg, int32_t code) {
1599
  SMqAskEpCbParam* pParam = (SMqAskEpCbParam*)param;
1600 1601 1602
  tmq_t*           tmq = taosAcquireRef(tmqMgmt.rsetId, pParam->refId);

  if (tmq == NULL) {
1603
    terrno = TSDB_CODE_TMQ_CONSUMER_CLOSED;
wmmhello's avatar
wmmhello 已提交
1604
//    pParam->pUserFn(tmq, terrno, NULL, pParam->pParam);
1605

1606
    taosMemoryFree(pMsg->pData);
dengyihao's avatar
dengyihao 已提交
1607
    taosMemoryFree(pMsg->pEpSet);
1608 1609
    taosMemoryFree(pParam);
    return terrno;
1610 1611
  }

H
Haojun Liao 已提交
1612
  if (code != TSDB_CODE_SUCCESS) {
1613 1614 1615 1616 1617 1618 1619 1620 1621
    tscError("consumer:0x%" PRIx64 ", get topic endpoint error, code:%s", tmq->consumerId, tstrerror(code));
    pParam->pUserFn(tmq, code, NULL, pParam->pParam);

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

    taosMemoryFree(pMsg->pData);
    taosMemoryFree(pMsg->pEpSet);
    taosMemoryFree(pParam);
    return code;
1622
  }
L
Liu Jicong 已提交
1623

L
Liu Jicong 已提交
1624
  // tmq's epoch is monotonically increase,
L
Liu Jicong 已提交
1625
  // so it's safe to discard any old epoch msg.
L
Liu Jicong 已提交
1626
  // Epoch will only increase when received newer epoch ep msg
L
Liu Jicong 已提交
1627 1628 1629
  SMqRspHead* head = pMsg->pData;
  int32_t     epoch = atomic_load_32(&tmq->epoch);
  if (head->epoch <= epoch) {
1630
    tscInfo("consumer:0x%" PRIx64 ", recv ep, msg epoch %d, current epoch %d, no need to update local ep",
1631
             tmq->consumerId, head->epoch, epoch);
1632

1633 1634 1635 1636 1637 1638 1639 1640
    if (tmq->status == TMQ_CONSUMER_STATUS__RECOVER) {
      SMqAskEpRsp rsp;
      tDecodeSMqAskEpRsp(POINTER_SHIFT(pMsg->pData, sizeof(SMqRspHead)), &rsp);
      int8_t flag = (taosArrayGetSize(rsp.topics) == 0) ? TMQ_CONSUMER_STATUS__NO_TOPIC : TMQ_CONSUMER_STATUS__READY;
      atomic_store_8(&tmq->status, flag);
      tDeleteSMqAskEpRsp(&rsp);
    }

X
Xiaoyu Wang 已提交
1641
  } else {
1642
    tscInfo("consumer:0x%" PRIx64 ", recv ep, msg epoch %d, current epoch %d, update local ep", tmq->consumerId,
1643
             head->epoch, epoch);
1644
  }
L
Liu Jicong 已提交
1645

1646
  pParam->pUserFn(tmq, code, pMsg, pParam->pParam);
1647 1648
  taosReleaseRef(tmqMgmt.rsetId, pParam->refId);

dengyihao's avatar
dengyihao 已提交
1649
  taosMemoryFree(pMsg->pEpSet);
L
Liu Jicong 已提交
1650
  taosMemoryFree(pMsg->pData);
1651
  taosMemoryFree(pParam);
L
Liu Jicong 已提交
1652
  return code;
1653 1654
}

L
Liu Jicong 已提交
1655
void tmqBuildConsumeReqImpl(SMqPollReq* pReq, tmq_t* tmq, int64_t timeout, SMqClientTopic* pTopic, SMqClientVg* pVg) {
L
Liu Jicong 已提交
1656 1657 1658 1659
  int32_t groupLen = strlen(tmq->groupId);
  memcpy(pReq->subKey, tmq->groupId, groupLen);
  pReq->subKey[groupLen] = TMQ_SEPARATOR;
  strcpy(pReq->subKey + groupLen + 1, pTopic->topicName);
1660

1661
  pReq->withTbName = tmq->withTbName;
L
Liu Jicong 已提交
1662
  pReq->consumerId = tmq->consumerId;
1663
  pReq->timeout = timeout;
X
Xiaoyu Wang 已提交
1664
  pReq->epoch = tmq->epoch;
H
Haojun Liao 已提交
1665
  pReq->reqOffset = pVg->offsetInfo.currentOffset;
D
dapan1121 已提交
1666
  pReq->head.vgId = pVg->vgId;
1667 1668
  pReq->useSnapshot = tmq->useSnapshot;
  pReq->reqId = generateRequestId();
1669 1670
}

L
Liu Jicong 已提交
1671 1672
SMqMetaRspObj* tmqBuildMetaRspFromWrapper(SMqPollRspWrapper* pWrapper) {
  SMqMetaRspObj* pRspObj = taosMemoryCalloc(1, sizeof(SMqMetaRspObj));
L
Liu Jicong 已提交
1673
  pRspObj->resType = RES_TYPE__TMQ_META;
L
Liu Jicong 已提交
1674 1675 1676 1677 1678 1679 1680 1681
  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;
}

1682
SMqRspObj* tmqBuildRspFromWrapper(SMqPollRspWrapper* pWrapper, SMqClientVg* pVg, int64_t* numOfRows) {
L
Liu Jicong 已提交
1683 1684
  SMqRspObj* pRspObj = taosMemoryCalloc(1, sizeof(SMqRspObj));
  pRspObj->resType = RES_TYPE__TMQ;
1685

1686
  (*numOfRows) = 0;
L
Liu Jicong 已提交
1687 1688
  tstrncpy(pRspObj->topic, pWrapper->topicHandle->topicName, TSDB_TOPIC_FNAME_LEN);
  tstrncpy(pRspObj->db, pWrapper->topicHandle->db, TSDB_DB_FNAME_LEN);
1689

L
Liu Jicong 已提交
1690
  pRspObj->vgId = pWrapper->vgHandle->vgId;
L
Liu Jicong 已提交
1691
  pRspObj->resIter = -1;
L
Liu Jicong 已提交
1692
  memcpy(&pRspObj->rsp, &pWrapper->dataRsp, sizeof(SMqDataRsp));
L
Liu Jicong 已提交
1693

L
Liu Jicong 已提交
1694 1695
  pRspObj->resInfo.totalRows = 0;
  pRspObj->resInfo.precision = TSDB_TIME_PRECISION_MILLI;
1696

L
Liu Jicong 已提交
1697
  if (!pWrapper->dataRsp.withSchema) {
L
Liu Jicong 已提交
1698 1699
    setResSchemaInfo(&pRspObj->resInfo, pWrapper->topicHandle->schema.pSchema, pWrapper->topicHandle->schema.nCols);
  }
L
Liu Jicong 已提交
1700

1701
  // extract the rows in this data packet
X
Xiaoyu Wang 已提交
1702
  for (int32_t i = 0; i < pRspObj->rsp.blockNum; ++i) {
1703
    SRetrieveTableRsp* pRetrieve = (SRetrieveTableRsp*)taosArrayGetP(pRspObj->rsp.blockData, i);
X
Xiaoyu Wang 已提交
1704
    int64_t            rows = htobe64(pRetrieve->numOfRows);
1705
    pVg->numOfRows += rows;
1706
    (*numOfRows) += rows;
1707 1708
  }

L
Liu Jicong 已提交
1709
  return pRspObj;
X
Xiaoyu Wang 已提交
1710 1711
}

1712
SMqTaosxRspObj* tmqBuildTaosxRspFromWrapper(SMqPollRspWrapper* pWrapper, SMqClientVg* pVg, int64_t* numOfRows) {
L
Liu Jicong 已提交
1713
  SMqTaosxRspObj* pRspObj = taosMemoryCalloc(1, sizeof(SMqTaosxRspObj));
1714
  pRspObj->resType = RES_TYPE__TMQ_METADATA;
L
Liu Jicong 已提交
1715 1716 1717 1718
  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;
1719
  memcpy(&pRspObj->rsp, &pWrapper->taosxRsp, sizeof(STaosxRsp));
L
Liu Jicong 已提交
1720 1721 1722

  pRspObj->resInfo.totalRows = 0;
  pRspObj->resInfo.precision = TSDB_TIME_PRECISION_MILLI;
1723
  if (!pWrapper->taosxRsp.withSchema) {
L
Liu Jicong 已提交
1724 1725 1726
    setResSchemaInfo(&pRspObj->resInfo, pWrapper->topicHandle->schema.pSchema, pWrapper->topicHandle->schema.nCols);
  }

1727 1728 1729 1730 1731 1732 1733
  // extract the rows in this data packet
  for (int32_t i = 0; i < pRspObj->rsp.blockNum; ++i) {
    SRetrieveTableRsp* pRetrieve = (SRetrieveTableRsp*)taosArrayGetP(pRspObj->rsp.blockData, i);
    int64_t            rows = htobe64(pRetrieve->numOfRows);
    pVg->numOfRows += rows;
    (*numOfRows) += rows;
  }
L
Liu Jicong 已提交
1734 1735 1736
  return pRspObj;
}

1737 1738 1739 1740 1741 1742 1743 1744 1745 1746 1747 1748 1749 1750 1751 1752 1753 1754 1755 1756 1757 1758 1759 1760 1761 1762 1763 1764 1765 1766 1767 1768 1769
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;
wmmhello's avatar
wmmhello 已提交
1770 1771 1772
//  pParam->pVg = pVg;  // pVg may be released,fix it
//  pParam->pTopic = pTopic;
  strcpy(pParam->topicName, pTopic->topicName);
1773
  pParam->vgId = pVg->vgId;
H
Haojun Liao 已提交
1774
  pParam->requestId = req.reqId;
1775 1776 1777 1778 1779 1780 1781 1782

  SMsgSendInfo* sendInfo = taosMemoryCalloc(1, sizeof(SMsgSendInfo));
  if (sendInfo == NULL) {
    taosMemoryFree(pParam);
    taosMemoryFree(msg);
    return handleErrorBeforePoll(pVg, pTmq);
  }

H
Haojun Liao 已提交
1783
  sendInfo->msgInfo = (SDataBuf){.pData = msg, .len = msgSize, .handle = NULL};
1784 1785 1786 1787 1788 1789 1790
  sendInfo->requestId = req.reqId;
  sendInfo->requestObjRefId = 0;
  sendInfo->param = pParam;
  sendInfo->fp = tmqPollCb;
  sendInfo->msgType = TDMT_VND_TMQ_CONSUME;

  int64_t transporterId = 0;
1791
  char    offsetFormatBuf[TSDB_OFFSET_LEN];
H
Haojun Liao 已提交
1792
  tFormatOffset(offsetFormatBuf, tListLen(offsetFormatBuf), &pVg->offsetInfo.currentOffset);
1793

X
Xiaoyu Wang 已提交
1794 1795
  tscDebug("consumer:0x%" PRIx64 " send poll to %s vgId:%d, epoch %d, req:%s, reqId:0x%" PRIx64, pTmq->consumerId,
           pTopic->topicName, pVg->vgId, pTmq->epoch, offsetFormatBuf, req.reqId);
1796 1797 1798
  asyncSendMsgToServer(pTmq->pTscObj->pAppInfo->pTransporter, &pVg->epSet, &transporterId, sendInfo);

  pVg->pollCnt++;
1799
  pVg->seekUpdated = false;   // reset this flag.
1800 1801 1802 1803 1804
  pTmq->pollCnt++;

  return TSDB_CODE_SUCCESS;
}

1805
// broadcast the poll request to all related vnodes
H
Haojun Liao 已提交
1806
static int32_t tmqPollImpl(tmq_t* tmq, int64_t timeout) {
1807 1808 1809
  if(atomic_load_8(&tmq->status) == TMQ_CONSUMER_STATUS__RECOVER){
    return 0;
  }
1810
  int32_t numOfTopics = taosArrayGetSize(tmq->clientTopics);
X
Xiaoyu Wang 已提交
1811
  tscDebug("consumer:0x%" PRIx64 " start to poll data, numOfTopics:%d", tmq->consumerId, numOfTopics);
1812 1813

  for (int i = 0; i < numOfTopics; i++) {
X
Xiaoyu Wang 已提交
1814
    SMqClientTopic* pTopic = taosArrayGet(tmq->clientTopics, i);
X
Xiaoyu Wang 已提交
1815
    int32_t         numOfVg = taosArrayGetSize(pTopic->vgs);
1816 1817

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

1825
      int32_t vgStatus = atomic_val_compare_exchange_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE, TMQ_VG_STATUS__WAIT);
1826
      if (vgStatus == TMQ_VG_STATUS__WAIT) {
L
Liu Jicong 已提交
1827
        int32_t vgSkipCnt = atomic_add_fetch_32(&pVg->vgSkipCnt, 1);
1828
        tscTrace("consumer:0x%" PRIx64 " epoch %d wait poll-rsp, skip vgId:%d skip cnt %d", tmq->consumerId, tmq->epoch,
X
Xiaoyu Wang 已提交
1829
                 pVg->vgId, vgSkipCnt);
X
Xiaoyu Wang 已提交
1830 1831
        continue;
      }
1832

L
Liu Jicong 已提交
1833
      atomic_store_32(&pVg->vgSkipCnt, 0);
1834 1835 1836
      int32_t code = doTmqPollImpl(tmq, pTopic, pVg, timeout);
      if (code != TSDB_CODE_SUCCESS) {
        return code;
D
dapan1121 已提交
1837
      }
X
Xiaoyu Wang 已提交
1838 1839
    }
  }
1840

1841
  tscDebug("consumer:0x%" PRIx64 " end to poll data", tmq->consumerId);
X
Xiaoyu Wang 已提交
1842 1843 1844
  return 0;
}

H
Haojun Liao 已提交
1845
static int32_t tmqHandleNoPollRsp(tmq_t* tmq, SMqRspWrapper* rspWrapper, bool* pReset) {
L
Liu Jicong 已提交
1846
  if (rspWrapper->tmqRspType == TMQ_MSG_TYPE__EP_RSP) {
L
fix  
Liu Jicong 已提交
1847
    /*printf("ep %d %d\n", rspMsg->head.epoch, tmq->epoch);*/
L
Liu Jicong 已提交
1848 1849
    if (rspWrapper->epoch > atomic_load_32(&tmq->epoch)) {
      SMqAskEpRspWrapper* pEpRspWrapper = (SMqAskEpRspWrapper*)rspWrapper;
L
Liu Jicong 已提交
1850
      SMqAskEpRsp*        rspMsg = &pEpRspWrapper->msg;
1851
      doUpdateLocalEp(tmq, rspWrapper->epoch, rspMsg);
L
temp  
Liu Jicong 已提交
1852
      /*tmqClearUnhandleMsg(tmq);*/
L
Liu Jicong 已提交
1853
      tDeleteSMqAskEpRsp(rspMsg);
X
Xiaoyu Wang 已提交
1854 1855
      *pReset = true;
    } else {
L
Liu Jicong 已提交
1856
      tmqFreeRspWrapper(rspWrapper);
X
Xiaoyu Wang 已提交
1857 1858 1859 1860 1861 1862 1863 1864
      *pReset = false;
    }
  } else {
    return -1;
  }
  return 0;
}

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

X
Xiaoyu Wang 已提交
1868
  while (1) {
1869 1870
    SMqRspWrapper* pRspWrapper = NULL;
    taosGetQitem(tmq->qall, (void**)&pRspWrapper);
1871

1872
    if (pRspWrapper == NULL) {
L
Liu Jicong 已提交
1873
      taosReadAllQitems(tmq->mqueue, tmq->qall);
1874 1875
      taosGetQitem(tmq->qall, (void**)&pRspWrapper);
      if (pRspWrapper == NULL) {
L
Liu Jicong 已提交
1876 1877
        return NULL;
      }
X
Xiaoyu Wang 已提交
1878 1879
    }

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

1882 1883
    if (pRspWrapper->tmqRspType == TMQ_MSG_TYPE__END_RSP) {
      taosFreeQitem(pRspWrapper);
L
Liu Jicong 已提交
1884
      terrno = TSDB_CODE_TQ_NO_COMMITTED_OFFSET;
H
Haojun Liao 已提交
1885
      tscError("consumer:0x%" PRIx64 " unexpected rsp from poll, code:%s", tmq->consumerId, tstrerror(terrno));
L
Liu Jicong 已提交
1886
      return NULL;
1887 1888
    } else if (pRspWrapper->tmqRspType == TMQ_MSG_TYPE__POLL_RSP) {
      SMqPollRspWrapper* pollRspWrapper = (SMqPollRspWrapper*)pRspWrapper;
H
Haojun Liao 已提交
1889

X
Xiaoyu Wang 已提交
1890
      int32_t     consumerEpoch = atomic_load_32(&tmq->epoch);
1891 1892 1893
      SMqDataRsp* pDataRsp = &pollRspWrapper->dataRsp;

      if (pDataRsp->head.epoch == consumerEpoch) {
wmmhello's avatar
wmmhello 已提交
1894 1895 1896 1897 1898 1899 1900 1901
        SMqClientVg* pVg = getVgInfo(tmq, pollRspWrapper->topicName, pollRspWrapper->vgId);
        pollRspWrapper->vgHandle = pVg;
        pollRspWrapper->topicHandle = getTopicInfo(tmq, pollRspWrapper->topicName);
        if(pollRspWrapper->vgHandle == NULL || pollRspWrapper->topicHandle == NULL){
          tscError("consumer:0x%" PRIx64 " get vg or topic error, topic:%s vgId:%d", tmq->consumerId,
                   pollRspWrapper->topicName, pollRspWrapper->vgId);
          return NULL;
        }
1902 1903 1904 1905
        // update the epset
        if (pollRspWrapper->pEpset != NULL) {
          SEp* pEp = GET_ACTIVE_EP(pollRspWrapper->pEpset);
          SEp* pOld = GET_ACTIVE_EP(&(pVg->epSet));
X
Xiaoyu Wang 已提交
1906 1907
          tscDebug("consumer:0x%" PRIx64 " update epset vgId:%d, ep:%s:%d, old ep:%s:%d", tmq->consumerId, pVg->vgId,
                   pEp->fqdn, pEp->port, pOld->fqdn, pOld->port);
1908 1909 1910
          pVg->epSet = *pollRspWrapper->pEpset;
        }

1911 1912 1913
        // update the local offset value only for the returned values, only when the local offset is NOT updated
        // by tmq_offset_seek function
        if (!pVg->seekUpdated) {
T
t_max 已提交
1914
          tscDebug("consumer:0x%" PRIx64" local offset is update, since seekupdate not set", tmq->consumerId);
1915
          pVg->offsetInfo.currentOffset = pDataRsp->rspOffset;
T
t_max 已提交
1916 1917
        } else {
          tscDebug("consumer:0x%" PRIx64" local offset is NOT update, since seekupdate is set", tmq->consumerId);
wmmhello's avatar
wmmhello 已提交
1918
        }
1919 1920

        // update the status
X
Xiaoyu Wang 已提交
1921
        atomic_store_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE);
H
Haojun Liao 已提交
1922

1923 1924 1925
        // update the valid wal version range
        pVg->offsetInfo.walVerBegin = pDataRsp->head.walsver;
        pVg->offsetInfo.walVerEnd = pDataRsp->head.walever;
1926
        pVg->receivedInfoFromVnode = true;
1927

1928 1929
        char buf[TSDB_OFFSET_LEN];
        tFormatOffset(buf, TSDB_OFFSET_LEN, &pDataRsp->rspOffset);
1930
        if (pDataRsp->blockNum == 0) {
X
Xiaoyu Wang 已提交
1931
          tscDebug("consumer:0x%" PRIx64 " empty block received, vgId:%d, offset:%s, vg total:%" PRId64
wmmhello's avatar
wmmhello 已提交
1932
                   ", total:%" PRId64 ", reqId:0x%" PRIx64,
X
Xiaoyu Wang 已提交
1933
                   tmq->consumerId, pVg->vgId, buf, pVg->numOfRows, tmq->totalRows, pollRspWrapper->reqId);
1934
          pRspWrapper = tmqFreeRspWrapper(pRspWrapper);
1935
          pVg->emptyBlockReceiveTs = taosGetTimestampMs();
L
Liu Jicong 已提交
1936
          taosFreeQitem(pollRspWrapper);
1937
        } else {  // build rsp
X
Xiaoyu Wang 已提交
1938
          int64_t    numOfRows = 0;
1939
          SMqRspObj* pRsp = tmqBuildRspFromWrapper(pollRspWrapper, pVg, &numOfRows);
H
Haojun Liao 已提交
1940
          tmq->totalRows += numOfRows;
1941
          pVg->emptyBlockReceiveTs = 0;
H
Haojun Liao 已提交
1942
          tscDebug("consumer:0x%" PRIx64 " process poll rsp, vgId:%d, offset:%s, blocks:%d, rows:%" PRId64
wmmhello's avatar
wmmhello 已提交
1943
                   ", vg total:%" PRId64 ", total:%" PRId64 ", reqId:0x%" PRIx64,
H
Haojun Liao 已提交
1944
                   tmq->consumerId, pVg->vgId, buf, pDataRsp->blockNum, numOfRows, pVg->numOfRows, tmq->totalRows,
H
Haojun Liao 已提交
1945
                   pollRspWrapper->reqId);
1946 1947 1948
          taosFreeQitem(pollRspWrapper);
          return pRsp;
        }
X
Xiaoyu Wang 已提交
1949
      } else {
H
Haojun Liao 已提交
1950
        tscDebug("consumer:0x%" PRIx64 " vgId:%d msg discard since epoch mismatch: msg epoch %d, consumer epoch %d",
1951
                 tmq->consumerId, pollRspWrapper->vgId, pDataRsp->head.epoch, consumerEpoch);
1952
        pRspWrapper = tmqFreeRspWrapper(pRspWrapper);
L
Liu Jicong 已提交
1953 1954
        taosFreeQitem(pollRspWrapper);
      }
1955
    } else if (pRspWrapper->tmqRspType == TMQ_MSG_TYPE__POLL_META_RSP) {
1956
      // todo handle the wal range and epset for each vgroup
1957
      SMqPollRspWrapper* pollRspWrapper = (SMqPollRspWrapper*)pRspWrapper;
L
Liu Jicong 已提交
1958
      int32_t            consumerEpoch = atomic_load_32(&tmq->epoch);
1959 1960 1961

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

L
Liu Jicong 已提交
1962
      if (pollRspWrapper->metaRsp.head.epoch == consumerEpoch) {
wmmhello's avatar
wmmhello 已提交
1963 1964 1965 1966 1967 1968 1969 1970
        SMqClientVg* pVg = getVgInfo(tmq, pollRspWrapper->topicName, pollRspWrapper->vgId);
        pollRspWrapper->vgHandle = pVg;
        pollRspWrapper->topicHandle = getTopicInfo(tmq, pollRspWrapper->topicName);
        if(pollRspWrapper->vgHandle == NULL || pollRspWrapper->topicHandle == NULL){
          tscError("consumer:0x%" PRIx64 " get vg or topic error, topic:%s vgId:%d", tmq->consumerId,
                   pollRspWrapper->topicName, pollRspWrapper->vgId);
          return NULL;
        }
H
Haojun Liao 已提交
1971

1972
        if(pollRspWrapper->metaRsp.rspOffset.type != 0){    // if offset is validate
H
Haojun Liao 已提交
1973
          pVg->offsetInfo.currentOffset = pollRspWrapper->metaRsp.rspOffset;
1974
        }
1975

L
Liu Jicong 已提交
1976 1977
        atomic_store_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE);
        // build rsp
L
Liu Jicong 已提交
1978
        SMqMetaRspObj* pRsp = tmqBuildMetaRspFromWrapper(pollRspWrapper);
L
Liu Jicong 已提交
1979 1980 1981
        taosFreeQitem(pollRspWrapper);
        return pRsp;
      } else {
H
Haojun Liao 已提交
1982
        tscDebug("consumer:0x%" PRIx64 " vgId:%d msg discard since epoch mismatch: msg epoch %d, consumer epoch %d",
1983
                 tmq->consumerId, pollRspWrapper->vgId, pollRspWrapper->metaRsp.head.epoch, consumerEpoch);
1984
        pRspWrapper = tmqFreeRspWrapper(pRspWrapper);
L
Liu Jicong 已提交
1985
        taosFreeQitem(pollRspWrapper);
X
Xiaoyu Wang 已提交
1986
      }
1987 1988
    } else if (pRspWrapper->tmqRspType == TMQ_MSG_TYPE__TAOSX_RSP) {
      SMqPollRspWrapper* pollRspWrapper = (SMqPollRspWrapper*)pRspWrapper;
X
Xiaoyu Wang 已提交
1989
      int32_t            consumerEpoch = atomic_load_32(&tmq->epoch);
H
Haojun Liao 已提交
1990

L
Liu Jicong 已提交
1991
      if (pollRspWrapper->taosxRsp.head.epoch == consumerEpoch) {
wmmhello's avatar
wmmhello 已提交
1992 1993 1994 1995 1996 1997 1998 1999
        SMqClientVg* pVg = getVgInfo(tmq, pollRspWrapper->topicName, pollRspWrapper->vgId);
        pollRspWrapper->vgHandle = pVg;
        pollRspWrapper->topicHandle = getTopicInfo(tmq, pollRspWrapper->topicName);
        if(pollRspWrapper->vgHandle == NULL || pollRspWrapper->topicHandle == NULL){
          tscError("consumer:0x%" PRIx64 " get vg or topic error, topic:%s vgId:%d", tmq->consumerId,
                   pollRspWrapper->topicName, pollRspWrapper->vgId);
          return NULL;
        }
H
Haojun Liao 已提交
2000

T
t_max 已提交
2001 2002 2003 2004 2005 2006 2007 2008 2009
        // update the local offset value only for the returned values, only when the local offset is NOT updated
        // by tmq_offset_seek function
        if (!pVg->seekUpdated) {
          if(pollRspWrapper->taosxRsp.rspOffset.type != 0) {    // if offset is validate
            tscDebug("consumer:0x%" PRIx64" local offset is update, since seekupdate not set", tmq->consumerId);
            pVg->offsetInfo.currentOffset = pollRspWrapper->taosxRsp.rspOffset;
          }
        } else {
          tscDebug("consumer:0x%" PRIx64" local offset is NOT update, since seekupdate is set", tmq->consumerId);
2010
        }
2011

L
Liu Jicong 已提交
2012
        atomic_store_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE);
H
Haojun Liao 已提交
2013

L
Liu Jicong 已提交
2014
        if (pollRspWrapper->taosxRsp.blockNum == 0) {
wmmhello's avatar
wmmhello 已提交
2015
          tscDebug("consumer:0x%" PRIx64 " taosx empty block received, vgId:%d, vg total:%" PRId64 ", reqId:0x%" PRIx64,
H
Haojun Liao 已提交
2016
                   tmq->consumerId, pVg->vgId, pVg->numOfRows, pollRspWrapper->reqId);
H
Haojun Liao 已提交
2017
          pVg->emptyBlockReceiveTs = taosGetTimestampMs();
2018
          pRspWrapper = tmqFreeRspWrapper(pRspWrapper);
H
Haojun Liao 已提交
2019
          taosFreeQitem(pollRspWrapper);
L
Liu Jicong 已提交
2020
          continue;
H
Haojun Liao 已提交
2021
        } else {
X
Xiaoyu Wang 已提交
2022
          pVg->emptyBlockReceiveTs = 0;  // reset the ts
L
Liu Jicong 已提交
2023
        }
wmmhello's avatar
wmmhello 已提交
2024

L
Liu Jicong 已提交
2025
        // build rsp
X
Xiaoyu Wang 已提交
2026
        void*   pRsp = NULL;
2027
        int64_t numOfRows = 0;
L
Liu Jicong 已提交
2028
        if (pollRspWrapper->taosxRsp.createTableNum == 0) {
2029
          pRsp = tmqBuildRspFromWrapper(pollRspWrapper, pVg, &numOfRows);
L
Liu Jicong 已提交
2030
        } else {
2031
          pRsp = tmqBuildTaosxRspFromWrapper(pollRspWrapper, pVg, &numOfRows);
wmmhello's avatar
wmmhello 已提交
2032
        }
H
Haojun Liao 已提交
2033

2034 2035
        tmq->totalRows += numOfRows;

2036 2037
        char buf[TSDB_OFFSET_LEN];
        tFormatOffset(buf, TSDB_OFFSET_LEN, &pVg->offsetInfo.currentOffset);
H
Haojun Liao 已提交
2038
        tscDebug("consumer:0x%" PRIx64 " process taosx poll rsp, vgId:%d, offset:%s, blocks:%d, rows:%" PRId64
wmmhello's avatar
wmmhello 已提交
2039
                 ", vg total:%" PRId64 ", total:%" PRId64 ", reqId:0x%" PRIx64,
H
Haojun Liao 已提交
2040
                 tmq->consumerId, pVg->vgId, buf, pollRspWrapper->dataRsp.blockNum, numOfRows, pVg->numOfRows,
H
Haojun Liao 已提交
2041
                 tmq->totalRows, pollRspWrapper->reqId);
H
Haojun Liao 已提交
2042 2043

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

L
Liu Jicong 已提交
2046
      } else {
H
Haojun Liao 已提交
2047
        tscDebug("consumer:0x%" PRIx64 " vgId:%d msg discard since epoch mismatch: msg epoch %d, consumer epoch %d",
2048
                 tmq->consumerId, pollRspWrapper->vgId, pollRspWrapper->taosxRsp.head.epoch, consumerEpoch);
2049
        pRspWrapper = tmqFreeRspWrapper(pRspWrapper);
L
Liu Jicong 已提交
2050 2051
        taosFreeQitem(pollRspWrapper);
      }
X
Xiaoyu Wang 已提交
2052
    } else {
H
Haojun Liao 已提交
2053 2054
      tscDebug("consumer:0x%" PRIx64 " not data msg received", tmq->consumerId);

X
Xiaoyu Wang 已提交
2055
      bool reset = false;
2056 2057
      tmqHandleNoPollRsp(tmq, pRspWrapper, &reset);
      taosFreeQitem(pRspWrapper);
X
Xiaoyu Wang 已提交
2058
      if (pollIfReset && reset) {
2059
        tscDebug("consumer:0x%" PRIx64 ", reset and repoll", tmq->consumerId);
2060
        tmqPollImpl(tmq, timeout);
X
Xiaoyu Wang 已提交
2061 2062 2063 2064 2065
      }
    }
  }
}

2066
TAOS_RES* tmq_consumer_poll(tmq_t* tmq, int64_t timeout) {
L
Liu Jicong 已提交
2067 2068
  void*   rspObj;
  int64_t startTime = taosGetTimestampMs();
L
Liu Jicong 已提交
2069

2070
  tscInfo("consumer:0x%" PRIx64 " start to poll at %" PRId64 ", timeout:%" PRId64, tmq->consumerId, startTime,
X
Xiaoyu Wang 已提交
2071
           timeout);
L
Liu Jicong 已提交
2072

2073
  // in no topic status, delayed task also need to be processed
L
Liu Jicong 已提交
2074
  if (atomic_load_8(&tmq->status) == TMQ_CONSUMER_STATUS__INIT) {
2075
    tscInfo("consumer:0x%" PRIx64 " poll return since consumer is init", tmq->consumerId);
2076
    taosMsleep(500);  //     sleep for a while
2077 2078 2079
    return NULL;
  }

wmmhello's avatar
wmmhello 已提交
2080
  while (atomic_load_8(&tmq->status) == TMQ_CONSUMER_STATUS__RECOVER) {
L
Liu Jicong 已提交
2081
    int32_t retryCnt = 0;
2082
    while (TSDB_CODE_MND_CONSUMER_NOT_READY == doAskEp(tmq)) {
H
Haojun Liao 已提交
2083
      if (retryCnt++ > 40) {
L
Liu Jicong 已提交
2084 2085
        return NULL;
      }
2086

2087
      tscInfo("consumer:0x%" PRIx64 " not ready, retry:%d/40 in 500ms", tmq->consumerId, retryCnt);
L
Liu Jicong 已提交
2088 2089 2090 2091
      taosMsleep(500);
    }
  }

X
Xiaoyu Wang 已提交
2092
  while (1) {
L
Liu Jicong 已提交
2093
    tmqHandleAllDelayedTask(tmq);
2094

L
Liu Jicong 已提交
2095
    if (tmqPollImpl(tmq, timeout) < 0) {
2096
      tscError("consumer:0x%" PRIx64 " return due to poll error", tmq->consumerId);
L
Liu Jicong 已提交
2097
    }
L
Liu Jicong 已提交
2098

2099
    rspObj = tmqHandleAllRsp(tmq, timeout, false);
L
Liu Jicong 已提交
2100
    if (rspObj) {
2101
      tscDebug("consumer:0x%" PRIx64 " return rsp %p", tmq->consumerId, rspObj);
L
Liu Jicong 已提交
2102
      return (TAOS_RES*)rspObj;
L
Liu Jicong 已提交
2103
    } else if (terrno == TSDB_CODE_TQ_NO_COMMITTED_OFFSET) {
2104
      tscInfo("consumer:0x%" PRIx64 " return null since no committed offset", tmq->consumerId);
L
Liu Jicong 已提交
2105
      return NULL;
X
Xiaoyu Wang 已提交
2106
    }
2107

2108
    if (timeout >= 0) {
L
Liu Jicong 已提交
2109
      int64_t currentTime = taosGetTimestampMs();
2110 2111
      int64_t elapsedTime = currentTime - startTime;
      if (elapsedTime > timeout) {
2112
        tscInfo("consumer:0x%" PRIx64 " (epoch %d) timeout, no rsp, start time %" PRId64 ", current time %" PRId64,
L
Liu Jicong 已提交
2113
                 tmq->consumerId, tmq->epoch, startTime, currentTime);
X
Xiaoyu Wang 已提交
2114 2115
        return NULL;
      }
2116
      tsem_timewait(&tmq->rspSem, (timeout - elapsedTime));
L
Liu Jicong 已提交
2117 2118
    } else {
      // use tsem_timewait instead of tsem_wait to avoid unexpected stuck
L
Liu Jicong 已提交
2119
      tsem_timewait(&tmq->rspSem, 1000);
X
Xiaoyu Wang 已提交
2120 2121 2122 2123
    }
  }
}

2124 2125 2126 2127 2128 2129 2130 2131 2132 2133 2134 2135 2136 2137
static void displayConsumeStatistics(const tmq_t* pTmq) {
  int32_t numOfTopics = taosArrayGetSize(pTmq->clientTopics);
  tscDebug("consumer:0x%" PRIx64 " closing poll:%" PRId64 " rows:%" PRId64 " topics:%d, final epoch:%d",
           pTmq->consumerId, pTmq->pollCnt, pTmq->totalRows, numOfTopics, pTmq->epoch);

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

    tscDebug("consumer:0x%" PRIx64 " topic:%d", pTmq->consumerId, i);
    int32_t numOfVgs = taosArrayGetSize(pTopics->vgs);
    for (int32_t j = 0; j < numOfVgs; ++j) {
      SMqClientVg* pVg = taosArrayGet(pTopics->vgs, j);
      tscDebug("topic:%s, %d. vgId:%d rows:%" PRId64, pTopics->topicName, j, pVg->vgId, pVg->numOfRows);
2138
    }
2139
  }
2140

2141 2142
  tscDebug("consumer:0x%" PRIx64 " rows dist end", pTmq->consumerId);
}
2143

2144
int32_t tmq_consumer_close(tmq_t* tmq) {
2145
  tscInfo("consumer:0x%" PRIx64 " start to close consumer, status:%d", tmq->consumerId, tmq->status);
2146
  displayConsumeStatistics(tmq);
2147

2148 2149 2150 2151 2152 2153
  if (tmq->status == TMQ_CONSUMER_STATUS__READY) {
    // if auto commit is set, commit before close consumer. Otherwise, do nothing.
    if (tmq->autoCommit) {
      int32_t rsp = tmq_commit_sync(tmq, NULL);
      if (rsp != 0) {
        return rsp;
2154 2155
      }
    }
2156
    taosSsleep(2);  // sleep 2s for hb to send offset and rows to server
2157

L
Liu Jicong 已提交
2158
    int32_t     retryCnt = 0;
2159
    tmq_list_t* lst = tmq_list_new();
L
Liu Jicong 已提交
2160
    while (1) {
2161
      int32_t rsp = tmq_subscribe(tmq, lst);
L
Liu Jicong 已提交
2162 2163 2164 2165 2166 2167 2168 2169
      if (rsp != TSDB_CODE_MND_CONSUMER_NOT_READY || retryCnt > 5) {
        break;
      } else {
        retryCnt++;
        taosMsleep(500);
      }
    }

2170
    tmq_list_destroy(lst);
2171
  } else {
2172
    tscInfo("consumer:0x%" PRIx64 " not in ready state, close it directly", tmq->consumerId);
L
Liu Jicong 已提交
2173
  }
H
Haojun Liao 已提交
2174

2175
  taosRemoveRef(tmqMgmt.rsetId, tmq->refId);
L
Liu Jicong 已提交
2176
  return 0;
2177
}
L
Liu Jicong 已提交
2178

L
Liu Jicong 已提交
2179 2180
const char* tmq_err2str(int32_t err) {
  if (err == 0) {
L
Liu Jicong 已提交
2181
    return "success";
L
Liu Jicong 已提交
2182
  } else if (err == -1) {
L
Liu Jicong 已提交
2183 2184 2185
    return "fail";
  } else {
    return tstrerror(err);
L
Liu Jicong 已提交
2186 2187
  }
}
L
Liu Jicong 已提交
2188

L
Liu Jicong 已提交
2189 2190 2191 2192 2193
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;
2194 2195
  } else if (TD_RES_TMQ_METADATA(res)) {
    return TMQ_RES_METADATA;
L
Liu Jicong 已提交
2196 2197 2198 2199 2200
  } else {
    return TMQ_RES_INVALID;
  }
}

L
Liu Jicong 已提交
2201
const char* tmq_get_topic_name(TAOS_RES* res) {
L
Liu Jicong 已提交
2202 2203
  if (TD_RES_TMQ(res)) {
    SMqRspObj* pRspObj = (SMqRspObj*)res;
L
Liu Jicong 已提交
2204
    return strchr(pRspObj->topic, '.') + 1;
L
Liu Jicong 已提交
2205 2206 2207
  } else if (TD_RES_TMQ_META(res)) {
    SMqMetaRspObj* pMetaRspObj = (SMqMetaRspObj*)res;
    return strchr(pMetaRspObj->topic, '.') + 1;
2208 2209 2210
  } else if (TD_RES_TMQ_METADATA(res)) {
    SMqTaosxRspObj* pRspObj = (SMqTaosxRspObj*)res;
    return strchr(pRspObj->topic, '.') + 1;
L
Liu Jicong 已提交
2211 2212 2213 2214 2215
  } else {
    return NULL;
  }
}

L
Liu Jicong 已提交
2216 2217 2218 2219
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 已提交
2220 2221 2222
  } else if (TD_RES_TMQ_META(res)) {
    SMqMetaRspObj* pMetaRspObj = (SMqMetaRspObj*)res;
    return strchr(pMetaRspObj->db, '.') + 1;
2223 2224 2225
  } else if (TD_RES_TMQ_METADATA(res)) {
    SMqTaosxRspObj* pRspObj = (SMqTaosxRspObj*)res;
    return strchr(pRspObj->db, '.') + 1;
L
Liu Jicong 已提交
2226 2227 2228 2229 2230
  } else {
    return NULL;
  }
}

L
Liu Jicong 已提交
2231 2232 2233 2234
int32_t tmq_get_vgroup_id(TAOS_RES* res) {
  if (TD_RES_TMQ(res)) {
    SMqRspObj* pRspObj = (SMqRspObj*)res;
    return pRspObj->vgId;
L
Liu Jicong 已提交
2235 2236 2237
  } else if (TD_RES_TMQ_META(res)) {
    SMqMetaRspObj* pMetaRspObj = (SMqMetaRspObj*)res;
    return pMetaRspObj->vgId;
2238
  } else if (TD_RES_TMQ_METADATA(res)) {
2239 2240
    SMqTaosxRspObj* pRspObj = (SMqTaosxRspObj*)res;
    return pRspObj->vgId;
L
Liu Jicong 已提交
2241 2242 2243 2244
  } else {
    return -1;
  }
}
L
Liu Jicong 已提交
2245

2246 2247 2248 2249 2250 2251 2252 2253 2254 2255 2256 2257 2258 2259 2260 2261 2262 2263 2264 2265 2266 2267 2268
int64_t tmq_get_vgroup_offset(TAOS_RES* res) {
  if (TD_RES_TMQ(res)) {
    SMqRspObj* pRspObj = (SMqRspObj*) res;
    STqOffsetVal* pOffset = &pRspObj->rsp.rspOffset;
    if (pOffset->type == TMQ_OFFSET__LOG) {
      return pRspObj->rsp.rspOffset.version;
    }
  } else if (TD_RES_TMQ_META(res)) {
    SMqMetaRspObj* pRspObj = (SMqMetaRspObj*)res;
    if (pRspObj->metaRsp.rspOffset.type == TMQ_OFFSET__LOG) {
      return pRspObj->metaRsp.rspOffset.version;
    }
  } else if (TD_RES_TMQ_METADATA(res)) {
    SMqTaosxRspObj* pRspObj = (SMqTaosxRspObj*) res;
    if (pRspObj->rsp.rspOffset.type == TMQ_OFFSET__LOG) {
      return pRspObj->rsp.rspOffset.version;
    }
  }

  // data from tsdb, no valid offset info
  return -1;
}

L
Liu Jicong 已提交
2269 2270 2271 2272 2273 2274 2275
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;
    }
2276
    return (const char*)taosArrayGetP(pRspObj->rsp.blockTbName, pRspObj->resIter);
2277 2278
  } else if (TD_RES_TMQ_METADATA(res)) {
    SMqTaosxRspObj* pRspObj = (SMqTaosxRspObj*)res;
L
Liu Jicong 已提交
2279 2280 2281
    if (!pRspObj->rsp.withTbName || pRspObj->rsp.blockTbName == NULL || pRspObj->resIter < 0 ||
        pRspObj->resIter >= pRspObj->rsp.blockNum) {
      return NULL;
2282
    }
L
Liu Jicong 已提交
2283 2284
    return (const char*)taosArrayGetP(pRspObj->rsp.blockTbName, pRspObj->resIter);
  }
L
Liu Jicong 已提交
2285 2286
  return NULL;
}
2287

2288 2289 2290 2291
void tmq_commit_async(tmq_t* tmq, const TAOS_RES* pRes, tmq_commit_cb* cb, void* param) {
  if (pRes == NULL) {  // here needs to commit all offsets.
    asyncCommitAllOffsets(tmq, cb, param);
  } else {  // only commit one offset
2292
    asyncCommitOffset(tmq, pRes, TDMT_VND_TMQ_COMMIT_OFFSET, cb, param);
2293
  }
L
Liu Jicong 已提交
2294 2295
}

2296
static void commitCallBackFn(tmq_t *UNUSED_PARAM(tmq), int32_t code, void* param) {
2297 2298 2299
  SSyncCommitInfo* pInfo = (SSyncCommitInfo*) param;
  pInfo->code = code;
  tsem_post(&pInfo->sem);
2300
}
2301

2302 2303 2304 2305 2306 2307 2308 2309 2310
int32_t tmq_commit_sync(tmq_t* tmq, const TAOS_RES* pRes) {
  int32_t code = 0;

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

  if (pRes == NULL) {
    asyncCommitAllOffsets(tmq, commitCallBackFn, pInfo);
H
Haojun Liao 已提交
2311
  } else {
2312
    asyncCommitOffset(tmq, pRes, TDMT_VND_TMQ_COMMIT_OFFSET, commitCallBackFn, pInfo);
2313 2314
  }

2315 2316
  tsem_wait(&pInfo->sem);
  code = pInfo->code;
H
Haojun Liao 已提交
2317 2318

  tsem_destroy(&pInfo->sem);
2319 2320
  taosMemoryFree(pInfo);

X
Xiaoyu Wang 已提交
2321
  tscDebug("consumer:0x%" PRIx64 " sync commit done, code:%s", tmq->consumerId, tstrerror(code));
2322 2323 2324 2325 2326 2327 2328 2329 2330 2331 2332 2333
  return code;
}

void updateEpCallbackFn(tmq_t* pTmq, int32_t code, SDataBuf* pDataBuf, void* param) {
  SAskEpInfo* pInfo = param;
  pInfo->code = code;

  if (code == TSDB_CODE_SUCCESS) {
    SMqRspHead* head = pDataBuf->pData;

    SMqAskEpRsp rsp;
    tDecodeSMqAskEpRsp(POINTER_SHIFT(pDataBuf->pData, sizeof(SMqRspHead)), &rsp);
2334
    doUpdateLocalEp(pTmq, head->epoch, &rsp);
2335 2336 2337
    tDeleteSMqAskEpRsp(&rsp);
  }

H
Haojun Liao 已提交
2338
  tsem_post(&pInfo->sem);
2339 2340 2341 2342 2343 2344 2345 2346 2347 2348 2349 2350 2351 2352 2353 2354 2355 2356 2357 2358 2359 2360 2361 2362 2363 2364
}

void addToQueueCallbackFn(tmq_t* pTmq, int32_t code, SDataBuf* pDataBuf, void* param) {
  if (code != TSDB_CODE_SUCCESS) {
    terrno = code;
    return;
  }

  SMqAskEpRspWrapper* pWrapper = taosAllocateQitem(sizeof(SMqAskEpRspWrapper), DEF_QITEM, 0);
  if (pWrapper == NULL) {
    terrno = TSDB_CODE_OUT_OF_MEMORY;
    return;
  }

  SMqRspHead* head = pDataBuf->pData;

  pWrapper->tmqRspType = TMQ_MSG_TYPE__EP_RSP;
  pWrapper->epoch = head->epoch;
  memcpy(&pWrapper->msg, pDataBuf->pData, sizeof(SMqRspHead));
  tDecodeSMqAskEpRsp(POINTER_SHIFT(pDataBuf->pData, sizeof(SMqRspHead)), &pWrapper->msg);

  taosWriteQitem(pTmq->mqueue, pWrapper);
}

int32_t doAskEp(tmq_t* pTmq) {
  SAskEpInfo* pInfo = taosMemoryMalloc(sizeof(SAskEpInfo));
H
Haojun Liao 已提交
2365
  tsem_init(&pInfo->sem, 0, 0);
2366 2367

  asyncAskEp(pTmq, updateEpCallbackFn, pInfo);
H
Haojun Liao 已提交
2368
  tsem_wait(&pInfo->sem);
2369 2370

  int32_t code = pInfo->code;
H
Haojun Liao 已提交
2371
  tsem_destroy(&pInfo->sem);
2372 2373 2374 2375 2376
  taosMemoryFree(pInfo);
  return code;
}

void asyncAskEp(tmq_t* pTmq, __tmq_askep_fn_t askEpFn, void* param) {
2377
  SMqAskEpReq req = {0};
2378 2379 2380
  req.consumerId = pTmq->consumerId;
  req.epoch = pTmq->epoch;
  strcpy(req.cgroup, pTmq->groupId);
2381 2382 2383

  int32_t tlen = tSerializeSMqAskEpReq(NULL, 0, &req);
  if (tlen < 0) {
2384 2385 2386
    tscError("consumer:0x%" PRIx64 ", tSerializeSMqAskEpReq failed", pTmq->consumerId);
    askEpFn(pTmq, TSDB_CODE_INVALID_PARA, NULL, param);
    return;
2387 2388 2389 2390
  }

  void* pReq = taosMemoryCalloc(1, tlen);
  if (pReq == NULL) {
2391 2392 2393
    tscError("consumer:0x%" PRIx64 ", failed to malloc askEpReq msg, size:%d", pTmq->consumerId, tlen);
    askEpFn(pTmq, TSDB_CODE_OUT_OF_MEMORY, NULL, param);
    return;
2394 2395 2396
  }

  if (tSerializeSMqAskEpReq(pReq, tlen, &req) < 0) {
2397
    tscError("consumer:0x%" PRIx64 ", tSerializeSMqAskEpReq %d failed", pTmq->consumerId, tlen);
2398
    taosMemoryFree(pReq);
2399 2400 2401

    askEpFn(pTmq, TSDB_CODE_INVALID_PARA, NULL, param);
    return;
2402 2403 2404 2405
  }

  SMqAskEpCbParam* pParam = taosMemoryCalloc(1, sizeof(SMqAskEpCbParam));
  if (pParam == NULL) {
2406
    tscError("consumer:0x%" PRIx64 ", failed to malloc subscribe param", pTmq->consumerId);
2407
    taosMemoryFree(pReq);
2408 2409 2410

    askEpFn(pTmq, TSDB_CODE_OUT_OF_MEMORY, NULL, param);
    return;
2411 2412
  }

2413 2414 2415 2416
  pParam->refId = pTmq->refId;
  pParam->epoch = pTmq->epoch;
  pParam->pUserFn = askEpFn;
  pParam->pParam = param;
2417 2418 2419 2420 2421

  SMsgSendInfo* sendInfo = taosMemoryCalloc(1, sizeof(SMsgSendInfo));
  if (sendInfo == NULL) {
    taosMemoryFree(pParam);
    taosMemoryFree(pReq);
2422 2423
    askEpFn(pTmq, TSDB_CODE_OUT_OF_MEMORY, NULL, param);
    return;
2424 2425
  }

X
Xiaoyu Wang 已提交
2426
  sendInfo->msgInfo = (SDataBuf){.pData = pReq, .len = tlen, .handle = NULL};
2427 2428 2429 2430

  sendInfo->requestId = generateRequestId();
  sendInfo->requestObjRefId = 0;
  sendInfo->param = pParam;
2431
  sendInfo->fp = askEpCallbackFn;
2432 2433
  sendInfo->msgType = TDMT_MND_TMQ_ASK_EP;

2434
  SEpSet epSet = getEpSet_s(&pTmq->pTscObj->pAppInfo->mgmtEp);
2435
  tscInfo("consumer:0x%" PRIx64 " ask ep from mnode, reqId:0x%" PRIx64, pTmq->consumerId, sendInfo->requestId);
2436 2437

  int64_t transporterId = 0;
2438
  asyncSendMsgToServer(pTmq->pTscObj->pAppInfo->pTransporter, &epSet, &transporterId, sendInfo);
2439 2440 2441 2442 2443 2444 2445
}

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

int32_t tmqCommitDone(SMqCommitCbParamSet* pParamSet) {
2446 2447 2448
  int64_t refId = pParamSet->refId;

  tmq_t* tmq = taosAcquireRef(tmqMgmt.rsetId, refId);
2449 2450 2451 2452 2453 2454 2455
  if (tmq == NULL) {
    taosMemoryFree(pParamSet);
    terrno = TSDB_CODE_TMQ_CONSUMER_CLOSED;
    return -1;
  }

  // if no more waiting rsp
2456
  pParamSet->callbackFn(tmq, pParamSet->code, pParamSet->userParam);
2457
  taosMemoryFree(pParamSet);
wmmhello's avatar
wmmhello 已提交
2458
//  tmq->needReportOffsetRows = true;
2459 2460

  taosReleaseRef(tmqMgmt.rsetId, refId);
2461
  return 0;
2462 2463
}

2464
void commitRspCountDown(SMqCommitCbParamSet* pParamSet, int64_t consumerId, const char* pTopic, int32_t vgId) {
2465 2466
  int32_t waitingRspNum = atomic_sub_fetch_32(&pParamSet->waitingRspNum, 1);
  if (waitingRspNum == 0) {
H
Haojun Liao 已提交
2467 2468
    tscDebug("consumer:0x%" PRIx64 " topic:%s vgId:%d all commit-rsp received, commit completed", consumerId, pTopic,
             vgId);
2469
    tmqCommitDone(pParamSet);
H
Haojun Liao 已提交
2470 2471 2472
  } else {
    tscDebug("consumer:0x%" PRIx64 " topic:%s vgId:%d commit-rsp received, remain:%d", consumerId, pTopic, vgId,
             waitingRspNum);
2473 2474
  }
}
2475 2476 2477 2478 2479 2480 2481 2482 2483 2484 2485 2486 2487 2488 2489 2490 2491 2492 2493 2494 2495 2496

SReqResultInfo* tmqGetNextResInfo(TAOS_RES* res, bool convertUcs4) {
  SMqRspObj* pRspObj = (SMqRspObj*)res;
  pRspObj->resIter++;

  if (pRspObj->resIter < pRspObj->rsp.blockNum) {
    SRetrieveTableRsp* pRetrieve = (SRetrieveTableRsp*)taosArrayGetP(pRspObj->rsp.blockData, pRspObj->resIter);
    if (pRspObj->rsp.withSchema) {
      SSchemaWrapper* pSW = (SSchemaWrapper*)taosArrayGetP(pRspObj->rsp.blockSchema, pRspObj->resIter);
      setResSchemaInfo(&pRspObj->resInfo, pSW->pSchema, pSW->nCols);
      taosMemoryFreeClear(pRspObj->resInfo.row);
      taosMemoryFreeClear(pRspObj->resInfo.pCol);
      taosMemoryFreeClear(pRspObj->resInfo.length);
      taosMemoryFreeClear(pRspObj->resInfo.convertBuf);
      taosMemoryFreeClear(pRspObj->resInfo.convertJson);
    }

    setQueryResultFromRsp(&pRspObj->resInfo, pRetrieve, convertUcs4, false);
    return &pRspObj->resInfo;
  }

  return NULL;
H
Haojun Liao 已提交
2497 2498
}

2499 2500 2501 2502 2503 2504 2505 2506 2507 2508 2509 2510 2511 2512 2513 2514 2515 2516 2517 2518 2519
static int32_t tmqGetWalInfoCb(void* param, SDataBuf* pMsg, int32_t code) {
  SMqVgWalInfoParam* pParam = param;
  SMqVgCommon* pCommon = pParam->pCommon;

  int32_t total = atomic_add_fetch_32(&pCommon->numOfRsp, 1);
  if (code != TSDB_CODE_SUCCESS) {
    tscError("consumer:0x%" PRIx64 " failed to get the wal info from vgId:%d for topic:%s", pCommon->consumerId,
             pParam->vgId, pCommon->pTopicName);
    pCommon->code = code;
  } else {
    SMqDataRsp rsp;
    SDecoder decoder;
    tDecoderInit(&decoder, POINTER_SHIFT(pMsg->pData, sizeof(SMqRspHead)), pMsg->len - sizeof(SMqRspHead));
    tDecodeMqDataRsp(&decoder, &rsp);
    tDecoderClear(&decoder);

    SMqRspHead* pHead = pMsg->pData;

    tmq_topic_assignment assignment = {.begin = pHead->walsver,
                                       .end = pHead->walever,
                                       .currentOffset = rsp.rspOffset.version,
2520
                                       .vgId = pParam->vgId};
2521 2522 2523 2524 2525 2526 2527 2528 2529 2530 2531 2532 2533 2534 2535 2536 2537 2538 2539 2540 2541 2542

    taosThreadMutexLock(&pCommon->mutex);
    taosArrayPush(pCommon->pList, &assignment);
    taosThreadMutexUnlock(&pCommon->mutex);
  }

  if (total == pParam->totalReq) {
    tsem_post(&pCommon->rsp);
  }

  taosMemoryFree(pParam);
  return 0;
}

static void destroyCommonInfo(SMqVgCommon* pCommon) {
  taosArrayDestroy(pCommon->pList);
  tsem_destroy(&pCommon->rsp);
  taosThreadMutexDestroy(&pCommon->mutex);
  taosMemoryFree(pCommon->pTopicName);
  taosMemoryFree(pCommon);
}

H
Haojun Liao 已提交
2543
int32_t tmq_get_topic_assignment(tmq_t* tmq, const char* pTopicName, tmq_topic_assignment** assignment,
H
Haojun Liao 已提交
2544
                                 int32_t* numOfAssignment) {
H
Haojun Liao 已提交
2545 2546 2547
  *numOfAssignment = 0;
  *assignment = NULL;

2548
  int32_t accId = tmq->pTscObj->acctId;
2549
  char    tname[128] = {0};
2550 2551 2552
  sprintf(tname, "%d.%s", accId, pTopicName);

  SMqClientTopic* pTopic = getTopicByName(tmq, tname);
H
Haojun Liao 已提交
2553 2554 2555 2556 2557 2558
  if (pTopic == NULL) {
    return TSDB_CODE_INVALID_PARA;
  }

  // in case of snapshot is opened, no valid offset will return
  *numOfAssignment = taosArrayGetSize(pTopic->vgs);
2559 2560 2561 2562 2563 2564 2565 2566

  *assignment = taosMemoryCalloc(*numOfAssignment, sizeof(tmq_topic_assignment));
  if (*assignment == NULL) {
    tscError("consumer:0x%" PRIx64 " failed to malloc buffer, size:%" PRIzu, tmq->consumerId,
             (*numOfAssignment) * sizeof(tmq_topic_assignment));
    return TSDB_CODE_OUT_OF_MEMORY;
  }

2567 2568
  bool needFetch = false;

H
Haojun Liao 已提交
2569 2570
  for (int32_t j = 0; j < (*numOfAssignment); ++j) {
    SMqClientVg* pClientVg = taosArrayGet(pTopic->vgs, j);
2571
    if (!pClientVg->receivedInfoFromVnode) {
2572 2573 2574
      needFetch = true;
      break;
    }
H
Haojun Liao 已提交
2575 2576 2577 2578 2579 2580 2581 2582 2583 2584

    tmq_topic_assignment* pAssignment = &(*assignment)[j];
    if (pClientVg->offsetInfo.currentOffset.type == TMQ_OFFSET__LOG) {
      pAssignment->currentOffset = pClientVg->offsetInfo.currentOffset.version;
    } else {
      pAssignment->currentOffset = 0;
    }

    pAssignment->begin = pClientVg->offsetInfo.walVerBegin;
    pAssignment->end = pClientVg->offsetInfo.walVerEnd;
2585
    pAssignment->vgId = pClientVg->vgId;
H
Haojun Liao 已提交
2586 2587
  }

2588 2589 2590 2591 2592 2593 2594 2595 2596 2597 2598 2599 2600 2601 2602 2603 2604 2605 2606 2607 2608 2609 2610 2611 2612 2613 2614 2615 2616 2617 2618 2619 2620 2621 2622 2623 2624 2625 2626 2627 2628 2629 2630 2631 2632 2633 2634 2635 2636 2637 2638 2639 2640 2641 2642 2643 2644 2645 2646 2647 2648 2649 2650 2651 2652 2653 2654 2655
  if (needFetch) {
    SMqVgCommon* pCommon = taosMemoryCalloc(1, sizeof(SMqVgCommon));
    if (pCommon == NULL) {
      terrno = TSDB_CODE_OUT_OF_MEMORY;
      return terrno;
    }

    pCommon->pList= taosArrayInit(4, sizeof(tmq_topic_assignment));
    tsem_init(&pCommon->rsp, 0, 0);
    taosThreadMutexInit(&pCommon->mutex, 0);
    pCommon->pTopicName = taosStrdup(pTopic->topicName);
    pCommon->consumerId = tmq->consumerId;

    terrno = TSDB_CODE_OUT_OF_MEMORY;
    for (int32_t i = 0; i < (*numOfAssignment); ++i) {
      SMqClientVg* pClientVg = taosArrayGet(pTopic->vgs, i);

      SMqVgWalInfoParam* pParam = taosMemoryMalloc(sizeof(SMqVgWalInfoParam));
      if (pParam == NULL) {
        destroyCommonInfo(pCommon);
        return terrno;
      }

      pParam->epoch = tmq->epoch;
      pParam->vgId = pClientVg->vgId;
      pParam->totalReq = *numOfAssignment;
      pParam->pCommon = pCommon;

      SMqPollReq req = {0};
      tmqBuildConsumeReqImpl(&req, tmq, 10, pTopic, pClientVg);

      int32_t msgSize = tSerializeSMqPollReq(NULL, 0, &req);
      if (msgSize < 0) {
        taosMemoryFree(pParam);
        destroyCommonInfo(pCommon);
        return terrno;
      }

      char* msg = taosMemoryCalloc(1, msgSize);
      if (NULL == msg) {
        taosMemoryFree(pParam);
        destroyCommonInfo(pCommon);
        return terrno;
      }

      if (tSerializeSMqPollReq(msg, msgSize, &req) < 0) {
        taosMemoryFree(msg);
        taosMemoryFree(pParam);
        destroyCommonInfo(pCommon);
        return terrno;
      }

      SMsgSendInfo* sendInfo = taosMemoryCalloc(1, sizeof(SMsgSendInfo));
      if (sendInfo == NULL) {
        taosMemoryFree(pParam);
        taosMemoryFree(msg);
        destroyCommonInfo(pCommon);
        return terrno;
      }

      sendInfo->msgInfo = (SDataBuf){.pData = msg, .len = msgSize, .handle = NULL};
      sendInfo->requestId = req.reqId;
      sendInfo->requestObjRefId = 0;
      sendInfo->param = pParam;
      sendInfo->fp = tmqGetWalInfoCb;
      sendInfo->msgType = TDMT_VND_TMQ_VG_WALINFO;

      int64_t transporterId = 0;
2656
      char    offsetFormatBuf[TSDB_OFFSET_LEN];
2657 2658
      tFormatOffset(offsetFormatBuf, tListLen(offsetFormatBuf), &pClientVg->offsetInfo.currentOffset);

2659
      tscInfo("consumer:0x%" PRIx64 " %s retrieve wal info vgId:%d, epoch %d, req:%s, reqId:0x%" PRIx64,
2660 2661 2662 2663 2664 2665 2666 2667 2668 2669
               tmq->consumerId, pTopic->topicName, pClientVg->vgId, tmq->epoch, offsetFormatBuf, req.reqId);
      asyncSendMsgToServer(tmq->pTscObj->pAppInfo->pTransporter, &pClientVg->epSet, &transporterId, sendInfo);
    }

    tsem_wait(&pCommon->rsp);
    int32_t code = pCommon->code;

    terrno = code;
    if (code != TSDB_CODE_SUCCESS) {
      taosMemoryFree(*assignment);
2670
      *assignment = NULL;
2671 2672 2673 2674 2675 2676 2677 2678 2679
      *numOfAssignment = 0;
    } else {
      int32_t num = taosArrayGetSize(pCommon->pList);
      for(int32_t i = 0; i < num; ++i) {
        (*assignment)[i] = *(tmq_topic_assignment*)taosArrayGet(pCommon->pList, i);
      }
      *numOfAssignment = num;
    }

2680 2681 2682 2683 2684 2685 2686 2687 2688 2689 2690 2691 2692
    for (int32_t j = 0; j < (*numOfAssignment); ++j) {
      tmq_topic_assignment* p = &(*assignment)[j];

      for(int32_t i = 0; i < taosArrayGetSize(pTopic->vgs); ++i) {
        SMqClientVg* pClientVg = taosArrayGet(pTopic->vgs, i);
        if (pClientVg->vgId != p->vgId) {
          continue;
        }

        SVgOffsetInfo* pOffsetInfo = &pClientVg->offsetInfo;

        pOffsetInfo->currentOffset.type = TMQ_OFFSET__LOG;

2693
        char offsetBuf[TSDB_OFFSET_LEN] = {0};
2694 2695
        tFormatOffset(offsetBuf, tListLen(offsetBuf), &pOffsetInfo->currentOffset);

2696
        tscInfo("vgId:%d offset is update to:%s", p->vgId, offsetBuf);
2697 2698 2699 2700 2701 2702 2703 2704

        pOffsetInfo->walVerBegin = p->begin;
        pOffsetInfo->walVerEnd = p->end;
        pOffsetInfo->currentOffset.version = p->currentOffset;
        pOffsetInfo->committedOffset.version = p->currentOffset;
      }
    }

2705 2706 2707 2708 2709
    destroyCommonInfo(pCommon);
    return code;
  } else {
    return TSDB_CODE_SUCCESS;
  }
H
Haojun Liao 已提交
2710 2711
}

T
t_max 已提交
2712 2713 2714 2715 2716 2717 2718 2719
void tmq_free_assignment(tmq_topic_assignment* pAssignment) {
    if (pAssignment == NULL) {
        return;
    }

    taosMemoryFree(pAssignment);
}

2720
int32_t tmq_offset_seek(tmq_t* tmq, const char* pTopicName, int32_t vgId, int64_t offset) {
H
Haojun Liao 已提交
2721
  if (tmq == NULL) {
H
Haojun Liao 已提交
2722
    tscError("invalid tmq handle, null");
H
Haojun Liao 已提交
2723 2724 2725
    return TSDB_CODE_INVALID_PARA;
  }

2726 2727 2728 2729 2730
  int32_t accId = tmq->pTscObj->acctId;
  char tname[128] = {0};
  sprintf(tname, "%d.%s", accId, pTopicName);

  SMqClientTopic* pTopic = getTopicByName(tmq, tname);
H
Haojun Liao 已提交
2731
  if (pTopic == NULL) {
2732
    tscError("consumer:0x%" PRIx64 " invalid topic name:%s", tmq->consumerId, pTopicName);
H
Haojun Liao 已提交
2733 2734 2735 2736
    return TSDB_CODE_INVALID_PARA;
  }

  SMqClientVg* pVg = NULL;
H
Haojun Liao 已提交
2737 2738
  int32_t      numOfVgs = taosArrayGetSize(pTopic->vgs);
  for (int32_t i = 0; i < numOfVgs; ++i) {
H
Haojun Liao 已提交
2739
    SMqClientVg* pClientVg = taosArrayGet(pTopic->vgs, i);
2740
    if (pClientVg->vgId == vgId) {
H
Haojun Liao 已提交
2741 2742 2743 2744 2745 2746
      pVg = pClientVg;
      break;
    }
  }

  if (pVg == NULL) {
2747
    tscError("consumer:0x%" PRIx64 " invalid vgroup id:%d", tmq->consumerId, vgId);
H
Haojun Liao 已提交
2748 2749 2750 2751 2752 2753
    return TSDB_CODE_INVALID_PARA;
  }

  SVgOffsetInfo* pOffsetInfo = &pVg->offsetInfo;

  int32_t type = pOffsetInfo->currentOffset.type;
2754
  if (type != TMQ_OFFSET__LOG && !OFFSET_IS_RESET_OFFSET(type)) {
2755
    tscError("consumer:0x%" PRIx64 " offset type:%d not wal version, seek not allowed", tmq->consumerId, type);
H
Haojun Liao 已提交
2756 2757 2758
    return TSDB_CODE_INVALID_PARA;
  }

2759
  if (type == TMQ_OFFSET__LOG && (offset < pOffsetInfo->walVerBegin || offset > pOffsetInfo->walVerEnd)) {
H
Haojun Liao 已提交
2760 2761
    tscError("consumer:0x%" PRIx64 " invalid seek params, offset:%" PRId64 ", valid range:[%" PRId64 ", %" PRId64 "]",
             tmq->consumerId, offset, pOffsetInfo->walVerBegin, pOffsetInfo->walVerEnd);
H
Haojun Liao 已提交
2762 2763 2764
    return TSDB_CODE_INVALID_PARA;
  }

H
Haojun Liao 已提交
2765 2766 2767
  // update the offset, and then commit to vnode
  if (pOffsetInfo->currentOffset.type == TMQ_OFFSET__LOG) {
    pOffsetInfo->currentOffset.version = offset;
2768
    pOffsetInfo->committedOffset.version = INT64_MIN;
2769
    pVg->seekUpdated = true;
H
Haojun Liao 已提交
2770 2771
  }

H
Haojun Liao 已提交
2772
  SMqRspObj rspObj = {.resType = RES_TYPE__TMQ, .vgId = pVg->vgId};
2773
  tstrncpy(rspObj.topic, tname, tListLen(rspObj.topic));
H
Haojun Liao 已提交
2774

2775
  tscInfo("consumer:0x%" PRIx64 " seek to %" PRId64 " on vgId:%d", tmq->consumerId, offset, pVg->vgId);
2776 2777 2778 2779 2780 2781 2782 2783 2784 2785 2786 2787 2788 2789 2790 2791 2792 2793

  SSyncCommitInfo* pInfo = taosMemoryMalloc(sizeof(SSyncCommitInfo));
  if (pInfo == NULL) {
    tscError("consumer:0x%"PRIx64" failed to prepare seek operation", tmq->consumerId);
    return TSDB_CODE_OUT_OF_MEMORY;
  }

  tsem_init(&pInfo->sem, 0, 0);
  pInfo->code = 0;

  asyncCommitOffset(tmq, &rspObj, TDMT_VND_TMQ_SEEK_TO_OFFSET, commitCallBackFn, pInfo);

  tsem_wait(&pInfo->sem);
  int32_t code = pInfo->code;

  tsem_destroy(&pInfo->sem);
  taosMemoryFree(pInfo);

H
Haojun Liao 已提交
2794 2795 2796 2797 2798 2799
  if (code != TSDB_CODE_SUCCESS) {
    tscError("consumer:0x%" PRIx64 " failed to send seek to vgId:%d, code:%s", tmq->consumerId, pVg->vgId,
             tstrerror(code));
  }

  return code;
2800
}