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

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

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;
143
  STqOffsetVal seekOffset;      // the first version in block for seek operation
H
Haojun Liao 已提交
144 145 146 147 148 149 150 151 152 153
  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 已提交
154
  int32_t       vgSkipCnt;              // here used to mark the slow vgroups
155
//  bool          receivedInfoFromVnode;  // has already received info from vnode
H
Haojun Liao 已提交
156 157
  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 已提交
158
  SEpSet        epSet;
159 160
} SMqClientVg;

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

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

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

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

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

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

218 219 220 221 222
typedef struct SMqSeekParam {
  tsem_t        sem;
  int32_t       code;
} SMqSeekParam;

223 224 225 226 227 228 229
typedef struct SMqVgWalInfoParam {
  int32_t      vgId;
  int32_t      epoch;
  int32_t      totalReq;
  SMqVgCommon* pCommon;
} SMqVgWalInfoParam;

230
typedef struct {
231 232
  int64_t        refId;
  int32_t        epoch;
L
Liu Jicong 已提交
233 234
  int32_t        waitingRspNum;
  int32_t        totalRspNum;
235
  int32_t        code;
236
  tmq_commit_cb* callbackFn;
L
Liu Jicong 已提交
237 238
  /*SArray*        successfulOffsets;*/
  /*SArray*        failedOffsets;*/
X
Xiaoyu Wang 已提交
239
  void* userParam;
240 241 242 243
} SMqCommitCbParamSet;

typedef struct {
  SMqCommitCbParamSet* params;
244
  SMqVgOffset*         pOffset;
H
Haojun Liao 已提交
245 246 247
  char                 topicName[TSDB_TOPIC_FNAME_LEN];
  int32_t              vgId;
  tmq_t*               pTmq;
248
} SMqCommitCbParam;
249

250 251 252 253 254
typedef struct SSyncCommitInfo {
  tsem_t  sem;
  int32_t code;
} SSyncCommitInfo;

255
static int32_t doAskEp(tmq_t* tmq);
256 257
static int32_t makeTopicVgroupKey(char* dst, const char* topicName, int32_t vg);
static int32_t tmqCommitDone(SMqCommitCbParamSet* pParamSet);
258
static int32_t doSendCommitMsg(tmq_t* tmq, SMqClientVg* pVg, const char* pTopicName, SMqCommitCbParamSet* pParamSet,
259
                               int32_t index, int32_t totalVgroups, int32_t type);
X
Xiaoyu Wang 已提交
260 261 262
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);
263

264
tmq_conf_t* tmq_conf_new() {
wafwerar's avatar
wafwerar 已提交
265
  tmq_conf_t* conf = taosMemoryCalloc(1, sizeof(tmq_conf_t));
266 267 268 269 270
  if (conf == NULL) {
    terrno = TSDB_CODE_OUT_OF_MEMORY;
    return conf;
  }

271
  conf->withTbName = false;
L
Liu Jicong 已提交
272
  conf->autoCommit = true;
273
  conf->autoCommitInterval = DEFAULT_AUTO_COMMIT_INTERVAL;
274
  conf->resetOffset = TMQ_OFFSET__RESET_EARLIEST;
275
  conf->hbBgEnable = true;
276

277 278 279
  return conf;
}

L
Liu Jicong 已提交
280
void tmq_conf_destroy(tmq_conf_t* conf) {
L
Liu Jicong 已提交
281
  if (conf) {
282 283 284 285 286 287 288 289 290
    if (conf->ip) {
      taosMemoryFree(conf->ip);
    }
    if (conf->user) {
      taosMemoryFree(conf->user);
    }
    if (conf->pass) {
      taosMemoryFree(conf->pass);
    }
L
Liu Jicong 已提交
291 292
    taosMemoryFree(conf);
  }
L
Liu Jicong 已提交
293 294 295
}

tmq_conf_res_t tmq_conf_set(tmq_conf_t* conf, const char* key, const char* value) {
296
  if (strcasecmp(key, "group.id") == 0) {
L
Liu Jicong 已提交
297
    tstrncpy(conf->groupId, value, TSDB_CGROUP_LEN);
L
Liu Jicong 已提交
298
    return TMQ_CONF_OK;
299
  }
L
Liu Jicong 已提交
300

301
  if (strcasecmp(key, "client.id") == 0) {
L
Liu Jicong 已提交
302
    tstrncpy(conf->clientId, value, 256);
L
Liu Jicong 已提交
303 304
    return TMQ_CONF_OK;
  }
L
Liu Jicong 已提交
305

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

318
  if (strcasecmp(key, "auto.commit.interval.ms") == 0) {
319
    conf->autoCommitInterval = taosStr2int64(value);
L
Liu Jicong 已提交
320 321 322
    return TMQ_CONF_OK;
  }

323 324 325
  if (strcasecmp(key, "auto.offset.reset") == 0) {
    if (strcasecmp(value, "none") == 0) {
      conf->resetOffset = TMQ_OFFSET__RESET_NONE;
L
Liu Jicong 已提交
326
      return TMQ_CONF_OK;
327
    } else if (strcasecmp(value, "earliest") == 0) {
328
      conf->resetOffset = TMQ_OFFSET__RESET_EARLIEST;
L
Liu Jicong 已提交
329
      return TMQ_CONF_OK;
330 331
    } else if (strcasecmp(value, "latest") == 0) {
      conf->resetOffset = TMQ_OFFSET__RESET_LATEST;
L
Liu Jicong 已提交
332 333 334 335 336
      return TMQ_CONF_OK;
    } else {
      return TMQ_CONF_INVALID;
    }
  }
L
Liu Jicong 已提交
337

338 339
  if (strcasecmp(key, "msg.with.table.name") == 0) {
    if (strcasecmp(value, "true") == 0) {
340
      conf->withTbName = true;
L
Liu Jicong 已提交
341
      return TMQ_CONF_OK;
342
    } else if (strcasecmp(value, "false") == 0) {
343
      conf->withTbName = false;
L
Liu Jicong 已提交
344
      return TMQ_CONF_OK;
345 346 347 348 349
    } else {
      return TMQ_CONF_INVALID;
    }
  }

350 351
  if (strcasecmp(key, "experimental.snapshot.enable") == 0) {
    if (strcasecmp(value, "true") == 0) {
L
Liu Jicong 已提交
352
      conf->snapEnable = true;
353
      return TMQ_CONF_OK;
354
    } else if (strcasecmp(value, "false") == 0) {
L
Liu Jicong 已提交
355
      conf->snapEnable = false;
356 357 358 359 360 361
      return TMQ_CONF_OK;
    } else {
      return TMQ_CONF_INVALID;
    }
  }

362
  if (strcasecmp(key, "experimental.snapshot.batch.size") == 0) {
363
    conf->snapBatchSize = taosStr2int64(value);
L
Liu Jicong 已提交
364 365 366
    return TMQ_CONF_OK;
  }

367
//  if (strcasecmp(key, "enable.heartbeat.background") == 0) {
X
Xiaoyu Wang 已提交
368 369 370 371 372 373 374
    //    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 {
375 376
//    tscError("the default value of enable.heartbeat.background is true, can not be seted");
//    return TMQ_CONF_INVALID;
X
Xiaoyu Wang 已提交
377
    //    }
378
//  }
L
Liu Jicong 已提交
379

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

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

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

395
  if (strcasecmp(key, "td.connect.port") == 0) {
396
    conf->port = taosStr2int64(value);
L
Liu Jicong 已提交
397 398
    return TMQ_CONF_OK;
  }
399

400
  if (strcasecmp(key, "td.connect.db") == 0) {
L
Liu Jicong 已提交
401 402 403
    return TMQ_CONF_OK;
  }

L
Liu Jicong 已提交
404
  return TMQ_CONF_UNKNOWN;
405 406
}

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

L
Liu Jicong 已提交
409 410
int32_t tmq_list_append(tmq_list_t* list, const char* src) {
  SArray* container = &list->container;
411
  if (src == NULL || src[0] == 0) return -1;
412
  char* topic = taosStrdup(src);
L
fix  
Liu Jicong 已提交
413
  if (taosArrayPush(container, &topic) == NULL) return -1;
414 415 416
  return 0;
}

L
Liu Jicong 已提交
417
void tmq_list_destroy(tmq_list_t* list) {
L
Liu Jicong 已提交
418
  SArray* container = &list->container;
L
Liu Jicong 已提交
419
  taosArrayDestroyP(container, taosMemoryFree);
L
Liu Jicong 已提交
420 421
}

L
Liu Jicong 已提交
422 423 424 425 426 427 428 429 430 431
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;
}

432 433 434 435 436 437 438 439 440 441 442 443 444 445 446 447 448 449 450 451 452 453 454 455
//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;
//}
456

457 458 459
// 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 已提交
460
static int32_t tmqCommitCb(void* param, SDataBuf* pBuf, int32_t code) {
461
  SMqCommitCbParam*    pParam = (SMqCommitCbParam*)param;
462
  SMqCommitCbParamSet* pParamSet = (SMqCommitCbParamSet*)pParam->params;
H
Haojun Liao 已提交
463

X
Xiaoyu Wang 已提交
464 465 466 467 468 469 470 471 472 473 474 475 476 477
  //  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 已提交
478
  //                 pParam->pTmq->consumerId, pParam->topicName, pVg->vgId, pVg->offsetInfo.committedOffset.version,
X
Xiaoyu Wang 已提交
479 480 481 482 483 484 485 486 487 488 489 490 491 492 493
  //                 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
494

L
Liu Jicong 已提交
495
  taosMemoryFree(pParam->pOffset);
L
Liu Jicong 已提交
496
  taosMemoryFree(pBuf->pData);
dengyihao's avatar
dengyihao 已提交
497
  taosMemoryFree(pBuf->pEpSet);
L
Liu Jicong 已提交
498

499
  commitRspCountDown(pParamSet, pParam->pTmq->consumerId, pParam->topicName, pParam->vgId);
500 501 502
  return 0;
}

503
static int32_t doSendCommitMsg(tmq_t* tmq, SMqClientVg* pVg, const char* pTopicName, SMqCommitCbParamSet* pParamSet,
504
                               int32_t index, int32_t totalVgroups, int32_t type) {
505
  SMqVgOffset* pOffset = taosMemoryCalloc(1, sizeof(SMqVgOffset));
L
Liu Jicong 已提交
506
  if (pOffset == NULL) {
507
    return TSDB_CODE_OUT_OF_MEMORY;
L
Liu Jicong 已提交
508
  }
509

510 511
  pOffset->consumerId = tmq->consumerId;
  pOffset->offset.val = pVg->offsetInfo.currentOffset;
512

L
Liu Jicong 已提交
513
  int32_t groupLen = strlen(tmq->groupId);
514 515 516
  memcpy(pOffset->offset.subKey, tmq->groupId, groupLen);
  pOffset->offset.subKey[groupLen] = TMQ_SEPARATOR;
  strcpy(pOffset->offset.subKey + groupLen + 1, pTopicName);
L
Liu Jicong 已提交
517

518 519
  int32_t len = 0;
  int32_t code = 0;
520
  tEncodeSize(tEncodeMqVgOffset, pOffset, len, code);
L
Liu Jicong 已提交
521
  if (code < 0) {
522
    return TSDB_CODE_INVALID_PARA;
L
Liu Jicong 已提交
523
  }
524

L
Liu Jicong 已提交
525
  void* buf = taosMemoryCalloc(1, sizeof(SMsgHead) + len);
L
Liu Jicong 已提交
526 527
  if (buf == NULL) {
    taosMemoryFree(pOffset);
528
    return TSDB_CODE_OUT_OF_MEMORY;
L
Liu Jicong 已提交
529
  }
530

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

L
Liu Jicong 已提交
533 534 535 536
  void* abuf = POINTER_SHIFT(buf, sizeof(SMsgHead));

  SEncoder encoder;
  tEncoderInit(&encoder, abuf, len);
537
  tEncodeMqVgOffset(&encoder, pOffset);
L
Liu Jicong 已提交
538
  tEncoderClear(&encoder);
L
Liu Jicong 已提交
539 540

  // build param
541
  SMqCommitCbParam* pParam = taosMemoryCalloc(1, sizeof(SMqCommitCbParam));
L
Liu Jicong 已提交
542
  if (pParam == NULL) {
L
Liu Jicong 已提交
543
    taosMemoryFree(pOffset);
L
Liu Jicong 已提交
544
    taosMemoryFree(buf);
545
    return TSDB_CODE_OUT_OF_MEMORY;
L
Liu Jicong 已提交
546
  }
547

L
Liu Jicong 已提交
548 549
  pParam->params = pParamSet;
  pParam->pOffset = pOffset;
H
Haojun Liao 已提交
550 551 552
  pParam->vgId = pVg->vgId;
  pParam->pTmq = tmq;

H
Haojun Liao 已提交
553
  tstrncpy(pParam->topicName, pTopicName, tListLen(pParam->topicName));
L
Liu Jicong 已提交
554 555 556 557

  // build send info
  SMsgSendInfo* pMsgSendInfo = taosMemoryCalloc(1, sizeof(SMsgSendInfo));
  if (pMsgSendInfo == NULL) {
L
Liu Jicong 已提交
558
    taosMemoryFree(pOffset);
L
Liu Jicong 已提交
559 560
    taosMemoryFree(buf);
    taosMemoryFree(pParam);
561
    return TSDB_CODE_OUT_OF_MEMORY;
L
Liu Jicong 已提交
562
  }
563

564
  pMsgSendInfo->msgInfo = (SDataBuf) { .pData = buf, .len = sizeof(SMsgHead) + len, .handle = NULL };
L
Liu Jicong 已提交
565 566 567 568

  pMsgSendInfo->requestId = generateRequestId();
  pMsgSendInfo->requestObjRefId = 0;
  pMsgSendInfo->param = pParam;
L
Liu Jicong 已提交
569
  pMsgSendInfo->paramFreeFp = taosMemoryFree;
570
  pMsgSendInfo->fp = tmqCommitCb;
571
  pMsgSendInfo->msgType = type;
L
Liu Jicong 已提交
572

L
Liu Jicong 已提交
573 574 575
  atomic_add_fetch_32(&pParamSet->waitingRspNum, 1);
  atomic_add_fetch_32(&pParamSet->totalRspNum, 1);

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

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

L
Liu Jicong 已提交
586 587
  int64_t transporterId = 0;
  asyncSendMsgToServer(tmq->pTscObj->pAppInfo->pTransporter, &pVg->epSet, &transporterId, pMsgSendInfo);
588 589

  return TSDB_CODE_SUCCESS;
L
Liu Jicong 已提交
590 591
}

H
Haojun Liao 已提交
592 593 594 595 596 597 598 599 600 601 602
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 已提交
603
  tscError("consumer:0x%" PRIx64 ", total:%d, failed to find topic:%s", tmq->consumerId, numOfTopics, pTopicName);
H
Haojun Liao 已提交
604 605 606
  return NULL;
}

607
static void asyncCommitOffset(tmq_t* tmq, const TAOS_RES* pRes, int32_t type, tmq_commit_cb* pCommitFp, void* userParam) {
608 609 610 611 612 613 614 615 616 617 618 619
  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 已提交
620
    vgId = pRspObj->vgId;
621 622 623
  } else if (TD_RES_TMQ_META(pRes)) {
    SMqMetaRspObj* pMetaRspObj = (SMqMetaRspObj*)pRes;
    pTopicName = pMetaRspObj->topic;
L
Liu Jicong 已提交
624
    vgId = pMetaRspObj->vgId;
625 626 627
  } else if (TD_RES_TMQ_METADATA(pRes)) {
    SMqTaosxRspObj* pRspObj = (SMqTaosxRspObj*)pRes;
    pTopicName = pRspObj->topic;
628
    vgId = pRspObj->vgId;
L
Liu Jicong 已提交
629
  } else {
630 631
    pCommitFp(tmq, TSDB_CODE_TMQ_INVALID_MSG, userParam);
    return;
L
Liu Jicong 已提交
632 633 634 635
  }

  SMqCommitCbParamSet* pParamSet = taosMemoryCalloc(1, sizeof(SMqCommitCbParamSet));
  if (pParamSet == NULL) {
636 637
    pCommitFp(tmq, TSDB_CODE_OUT_OF_MEMORY, userParam);
    return;
L
Liu Jicong 已提交
638
  }
H
Haojun Liao 已提交
639

640 641
  pParamSet->refId = tmq->refId;
  pParamSet->epoch = tmq->epoch;
642
  pParamSet->callbackFn = pCommitFp;
L
Liu Jicong 已提交
643
  pParamSet->userParam = userParam;
L
Liu Jicong 已提交
644

wmmhello's avatar
wmmhello 已提交
645
  taosRLockLatch(&tmq->lock);
H
Haojun Liao 已提交
646 647
  int32_t numOfTopics = taosArrayGetSize(tmq->clientTopics);

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

H
Haojun Liao 已提交
650 651
  SMqClientTopic* pTopic = getTopicByName(tmq, pTopicName);
  if (pTopic == NULL) {
X
Xiaoyu Wang 已提交
652 653
    tscWarn("consumer:0x%" PRIx64 " failed to find the specified topic:%s, total topics:%d", tmq->consumerId,
            pTopicName, numOfTopics);
654 655
    taosMemoryFree(pParamSet);
    pCommitFp(tmq, TSDB_CODE_SUCCESS, userParam);
wmmhello's avatar
wmmhello 已提交
656
    taosRUnLockLatch(&tmq->lock);
657 658
    return;
  }
L
Liu Jicong 已提交
659

660 661 662
  int32_t j = 0;
  int32_t numOfVgroups = taosArrayGetSize(pTopic->vgs);
  for (j = 0; j < numOfVgroups; j++) {
wmmhello's avatar
wmmhello 已提交
663
    SMqClientVg* pVg = (SMqClientVg*)taosArrayGet(pTopic->vgs, j);
664 665
    if (pVg->vgId == vgId) {
      break;
L
Liu Jicong 已提交
666
    }
L
Liu Jicong 已提交
667
  }
L
Liu Jicong 已提交
668

669
  if (j == numOfVgroups) {
X
Xiaoyu Wang 已提交
670 671
    tscWarn("consumer:0x%" PRIx64 " failed to find the specified vgId:%d, total Vgs:%d, topic:%s", tmq->consumerId,
            vgId, numOfVgroups, pTopicName);
L
Liu Jicong 已提交
672
    taosMemoryFree(pParamSet);
673
    pCommitFp(tmq, TSDB_CODE_SUCCESS, userParam);
wmmhello's avatar
wmmhello 已提交
674
    taosRUnLockLatch(&tmq->lock);
675
    return;
L
Liu Jicong 已提交
676 677
  }

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

682 683 684 685 686
    // failed to commit, callback user function directly.
    if (code != TSDB_CODE_SUCCESS) {
      taosMemoryFree(pParamSet);
      pCommitFp(tmq, code, userParam);
    }
687 688
    // update the offset value.
    pVg->offsetInfo.committedOffset = pVg->offsetInfo.currentOffset;
X
Xiaoyu Wang 已提交
689
  } else {  // do not perform commit, callback user function directly.
690 691
    taosMemoryFree(pParamSet);
    pCommitFp(tmq, code, userParam);
L
Liu Jicong 已提交
692
  }
wmmhello's avatar
wmmhello 已提交
693
  taosRUnLockLatch(&tmq->lock);
L
Liu Jicong 已提交
694 695
}

696
static void asyncCommitAllOffsets(tmq_t* tmq, tmq_commit_cb* pCommitFp, void* userParam) {
697 698
  SMqCommitCbParamSet* pParamSet = taosMemoryCalloc(1, sizeof(SMqCommitCbParamSet));
  if (pParamSet == NULL) {
699 700
    pCommitFp(tmq, TSDB_CODE_OUT_OF_MEMORY, userParam);
    return;
701
  }
702 703 704

  pParamSet->refId = tmq->refId;
  pParamSet->epoch = tmq->epoch;
705
  pParamSet->callbackFn = pCommitFp;
706 707
  pParamSet->userParam = userParam;

708 709 710
  // init as 1 to prevent concurrency issue
  pParamSet->waitingRspNum = 1;

wmmhello's avatar
wmmhello 已提交
711
  taosRLockLatch(&tmq->lock);
712
  int32_t numOfTopics = taosArrayGetSize(tmq->clientTopics);
X
Xiaoyu Wang 已提交
713
  tscDebug("consumer:0x%" PRIx64 " start to commit offset for %d topics", tmq->consumerId, numOfTopics);
714 715

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

719 720
    tscDebug("consumer:0x%" PRIx64 " commit offset for topics:%s, numOfVgs:%d", tmq->consumerId, pTopic->topicName,
             numOfVgroups);
721
    for (int32_t j = 0; j < numOfVgroups; j++) {
722 723
      SMqClientVg* pVg = taosArrayGet(pTopic->vgs, j);

H
Haojun Liao 已提交
724
      if (pVg->offsetInfo.currentOffset.type > 0 && !tOffsetEqual(&pVg->offsetInfo.currentOffset, &pVg->offsetInfo.committedOffset)) {
725
        int32_t code = doSendCommitMsg(tmq, pVg, pTopic->topicName, pParamSet, j, numOfVgroups, TDMT_VND_TMQ_COMMIT_OFFSET);
726 727
        if (code != TSDB_CODE_SUCCESS) {
          tscError("consumer:0x%" PRIx64 " topic:%s vgId:%d offset:%" PRId64 " failed, code:%s ordinal:%d/%d",
H
Haojun Liao 已提交
728
                   tmq->consumerId, pTopic->topicName, pVg->vgId, pVg->offsetInfo.committedOffset.version, tstrerror(terrno),
729
                   j + 1, numOfVgroups);
L
Liu Jicong 已提交
730 731
          continue;
        }
H
Haojun Liao 已提交
732 733

        // update the offset value.
H
Haojun Liao 已提交
734
        pVg->offsetInfo.committedOffset = pVg->offsetInfo.currentOffset;
735
      } else {
D
dapan1121 已提交
736
        tscDebug("consumer:0x%" PRIx64 " topic:%s vgId:%d, no commit, current:%" PRId64 ", ordinal:%d/%d",
H
Haojun Liao 已提交
737
                 tmq->consumerId, pTopic->topicName, pVg->vgId, pVg->offsetInfo.currentOffset.version, j + 1, numOfVgroups);
738 739 740
      }
    }
  }
wmmhello's avatar
wmmhello 已提交
741
  taosRUnLockLatch(&tmq->lock);
742

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

L
Liu Jicong 已提交
746
  // no request is sent
L
Liu Jicong 已提交
747 748
  if (pParamSet->totalRspNum == 0) {
    taosMemoryFree(pParamSet);
749 750
    pCommitFp(tmq, TSDB_CODE_SUCCESS, userParam);
    return;
L
Liu Jicong 已提交
751 752
  }

L
Liu Jicong 已提交
753
  // count down since waiting rsp num init as 1
754
  commitRspCountDown(pParamSet, tmq->consumerId, "", 0);
755 756
}

757 758
static void generateTimedTask(int64_t refId, int32_t type) {
  tmq_t* tmq = taosAcquireRef(tmqMgmt.rsetId, refId);
wmmhello's avatar
wmmhello 已提交
759 760 761 762 763 764 765 766 767
  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);
768 769 770 771 772
}

void tmqAssignAskEpTask(void* param, void* tmrId) {
  int64_t refId = *(int64_t*)param;
  generateTimedTask(refId, TMQ_DELAYED_TASK__ASK_EP);
773
  taosMemoryFree(param);
L
Liu Jicong 已提交
774 775 776
}

void tmqAssignDelayedCommitTask(void* param, void* tmrId) {
777
  int64_t refId = *(int64_t*)param;
778
  generateTimedTask(refId, TMQ_DELAYED_TASK__COMMIT);
779
  taosMemoryFree(param);
L
Liu Jicong 已提交
780 781
}

wmmhello's avatar
wmmhello 已提交
782 783 784 785 786 787 788 789 790 791 792 793 794
//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 已提交
795

796
int32_t tmqHbCb(void* param, SDataBuf* pMsg, int32_t code) {
dengyihao's avatar
dengyihao 已提交
797 798 799 800
  if (pMsg) {
    taosMemoryFree(pMsg->pData);
    taosMemoryFree(pMsg->pEpSet);
  }
801 802 803 804
  return 0;
}

void tmqSendHbReq(void* param, void* tmrId) {
805
  int64_t refId = *(int64_t*)param;
806

X
Xiaoyu Wang 已提交
807
  tmq_t* tmq = taosAcquireRef(tmqMgmt.rsetId, refId);
808
  if (tmq == NULL) {
L
Liu Jicong 已提交
809
    taosMemoryFree(param);
810 811
    return;
  }
D
dapan1121 已提交
812 813 814 815

  SMqHbReq req = {0};
  req.consumerId = tmq->consumerId;
  req.epoch = tmq->epoch;
wmmhello's avatar
wmmhello 已提交
816
  taosRLockLatch(&tmq->lock);
wmmhello's avatar
wmmhello 已提交
817
//  if(tmq->needReportOffsetRows){
818 819 820 821 822 823 824 825 826 827 828 829
    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;
830
        offRows->offset = pVg->offsetInfo.seekOffset;
831 832
        char buf[TSDB_OFFSET_LEN] = {0};
        tFormatOffset(buf, TSDB_OFFSET_LEN, &offRows->offset);
wmmhello's avatar
wmmhello 已提交
833
        tscInfo("consumer:0x%" PRIx64 ",report offset: vgId:%d, offset:%s, rows:%"PRId64, tmq->consumerId, offRows->vgId, buf, offRows->rows);
834
      }
835
    }
wmmhello's avatar
wmmhello 已提交
836 837
//    tmq->needReportOffsetRows = false;
//  }
wmmhello's avatar
wmmhello 已提交
838
  taosRUnLockLatch(&tmq->lock);
D
dapan1121 已提交
839

L
Liu Jicong 已提交
840
  int32_t tlen = tSerializeSMqHbReq(NULL, 0, &req);
D
dapan1121 已提交
841 842
  if (tlen < 0) {
    tscError("tSerializeSMqHbReq failed");
843
    goto OVER;
D
dapan1121 已提交
844
  }
845

L
Liu Jicong 已提交
846
  void* pReq = taosMemoryCalloc(1, tlen);
D
dapan1121 已提交
847 848
  if (tlen < 0) {
    tscError("failed to malloc MqHbReq msg, size:%d", tlen);
849
    goto OVER;
D
dapan1121 已提交
850
  }
851

D
dapan1121 已提交
852 853 854
  if (tSerializeSMqHbReq(pReq, tlen, &req) < 0) {
    tscError("tSerializeSMqHbReq %d failed", tlen);
    taosMemoryFree(pReq);
855
    goto OVER;
D
dapan1121 已提交
856
  }
857 858 859 860

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

864
  sendInfo->msgInfo = (SDataBuf){ .pData = pReq, .len = tlen, .handle = NULL };
865 866 867 868 869

  sendInfo->requestId = generateRequestId();
  sendInfo->requestObjRefId = 0;
  sendInfo->param = NULL;
  sendInfo->fp = tmqHbCb;
L
Liu Jicong 已提交
870
  sendInfo->msgType = TDMT_MND_TMQ_HB;
871 872 873 874 875 876 877

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

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

OVER:
878
  tDeatroySMqHbReq(&req);
879
  taosTmrReset(tmqSendHbReq, 1000, param, tmqMgmt.timer, &tmq->hbLiveTimer);
880
  taosReleaseRef(tmqMgmt.rsetId, refId);
881 882
}

883 884
static void defaultCommitCbFn(tmq_t* pTmq, int32_t code, void* param) {
  if (code != 0) {
885
    tscError("consumer:0x%" PRIx64 ", failed to commit offset, code:%s", pTmq->consumerId, tstrerror(code));
886 887 888
  }
}

889
int32_t tmqHandleAllDelayedTask(tmq_t* pTmq) {
L
Liu Jicong 已提交
890
  STaosQall* qall = taosAllocateQall();
891
  taosReadAllQitems(pTmq->delayedTask, qall);
L
Liu Jicong 已提交
892

893 894 895 896
  if (qall->numOfItems == 0) {
    taosFreeQall(qall);
    return TSDB_CODE_SUCCESS;
  }
897

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

902
  while (pTaskType != NULL) {
903
    if (*pTaskType == TMQ_DELAYED_TASK__ASK_EP) {
904
      asyncAskEp(pTmq, addToQueueCallbackFn, NULL);
905 906

      int64_t* pRefId = taosMemoryMalloc(sizeof(int64_t));
907
      *pRefId = pTmq->refId;
908

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

      asyncCommitAllOffsets(pTmq, pCallbackFn, pTmq->commitCbUserParam);
915
      int64_t* pRefId = taosMemoryMalloc(sizeof(int64_t));
916
      *pRefId = pTmq->refId;
917

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

L
Liu Jicong 已提交
924
    taosFreeQitem(pTaskType);
925
    taosGetQitem(qall, (void**)&pTaskType);
L
Liu Jicong 已提交
926
  }
927

L
Liu Jicong 已提交
928 929 930 931
  taosFreeQall(qall);
  return 0;
}

932
static void* tmqFreeRspWrapper(SMqRspWrapper* rspWrapper) {
L
Liu Jicong 已提交
933 934 935 936 937 938 939
  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;
940 941
    taosMemoryFreeClear(pRsp->pEpset);

L
Liu Jicong 已提交
942 943 944
    taosArrayDestroyP(pRsp->dataRsp.blockData, taosMemoryFree);
    taosArrayDestroy(pRsp->dataRsp.blockDataLen);
    taosArrayDestroyP(pRsp->dataRsp.blockTbName, taosMemoryFree);
945
    taosArrayDestroyP(pRsp->dataRsp.blockSchema, (FDelete)tDeleteSchemaWrapper);
L
Liu Jicong 已提交
946 947
  } else if (rspWrapper->tmqRspType == TMQ_MSG_TYPE__POLL_META_RSP) {
    SMqPollRspWrapper* pRsp = (SMqPollRspWrapper*)rspWrapper;
948 949
    taosMemoryFreeClear(pRsp->pEpset);

L
Liu Jicong 已提交
950 951 952
    taosMemoryFree(pRsp->metaRsp.metaRsp);
  } else if (rspWrapper->tmqRspType == TMQ_MSG_TYPE__TAOSX_RSP) {
    SMqPollRspWrapper* pRsp = (SMqPollRspWrapper*)rspWrapper;
953 954
    taosMemoryFreeClear(pRsp->pEpset);

L
Liu Jicong 已提交
955 956 957
    taosArrayDestroyP(pRsp->taosxRsp.blockData, taosMemoryFree);
    taosArrayDestroy(pRsp->taosxRsp.blockDataLen);
    taosArrayDestroyP(pRsp->taosxRsp.blockTbName, taosMemoryFree);
958
    taosArrayDestroyP(pRsp->taosxRsp.blockSchema, (FDelete)tDeleteSchemaWrapper);
L
Liu Jicong 已提交
959 960 961 962
    // taosx
    taosArrayDestroy(pRsp->taosxRsp.createTableLen);
    taosArrayDestroyP(pRsp->taosxRsp.createTableReq, taosMemoryFree);
  }
963 964

  return NULL;
L
Liu Jicong 已提交
965 966
}

L
Liu Jicong 已提交
967
void tmqClearUnhandleMsg(tmq_t* tmq) {
L
Liu Jicong 已提交
968
  SMqRspWrapper* rspWrapper = NULL;
L
Liu Jicong 已提交
969
  while (1) {
L
Liu Jicong 已提交
970 971 972 973 974
    taosGetQitem(tmq->qall, (void**)&rspWrapper);
    if (rspWrapper) {
      tmqFreeRspWrapper(rspWrapper);
      taosFreeQitem(rspWrapper);
    } else {
L
Liu Jicong 已提交
975
      break;
L
Liu Jicong 已提交
976
    }
L
Liu Jicong 已提交
977 978
  }

L
Liu Jicong 已提交
979
  rspWrapper = NULL;
L
Liu Jicong 已提交
980 981
  taosReadAllQitems(tmq->mqueue, tmq->qall);
  while (1) {
L
Liu Jicong 已提交
982 983 984 985 986
    taosGetQitem(tmq->qall, (void**)&rspWrapper);
    if (rspWrapper) {
      tmqFreeRspWrapper(rspWrapper);
      taosFreeQitem(rspWrapper);
    } else {
L
Liu Jicong 已提交
987
      break;
L
Liu Jicong 已提交
988
    }
L
Liu Jicong 已提交
989 990 991
  }
}

D
dapan1121 已提交
992
int32_t tmqSubscribeCb(void* param, SDataBuf* pMsg, int32_t code) {
L
Liu Jicong 已提交
993 994
  SMqSubscribeCbParam* pParam = (SMqSubscribeCbParam*)param;
  pParam->rspErr = code;
dengyihao's avatar
dengyihao 已提交
995 996

  taosMemoryFree(pMsg->pEpSet);
L
Liu Jicong 已提交
997 998 999
  tsem_post(&pParam->rspSem);
  return 0;
}
1000

L
Liu Jicong 已提交
1001
int32_t tmq_subscription(tmq_t* tmq, tmq_list_t** topics) {
X
Xiaoyu Wang 已提交
1002 1003 1004
  if (*topics == NULL) {
    *topics = tmq_list_new();
  }
wmmhello's avatar
wmmhello 已提交
1005
  taosRLockLatch(&tmq->lock);
X
Xiaoyu Wang 已提交
1006
  for (int i = 0; i < taosArrayGetSize(tmq->clientTopics); i++) {
L
Liu Jicong 已提交
1007
    SMqClientTopic* topic = taosArrayGet(tmq->clientTopics, i);
L
Liu Jicong 已提交
1008
    tmq_list_append(*topics, strchr(topic->topicName, '.') + 1);
X
Xiaoyu Wang 已提交
1009
  }
wmmhello's avatar
wmmhello 已提交
1010
  taosRUnLockLatch(&tmq->lock);
L
Liu Jicong 已提交
1011
  return 0;
X
Xiaoyu Wang 已提交
1012 1013
}

L
Liu Jicong 已提交
1014
int32_t tmq_unsubscribe(tmq_t* tmq) {
1015 1016 1017 1018 1019 1020 1021 1022
  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 已提交
1023 1024
  int32_t     rsp;
  int32_t     retryCnt = 0;
L
Liu Jicong 已提交
1025
  tmq_list_t* lst = tmq_list_new();
L
Liu Jicong 已提交
1026 1027 1028 1029 1030 1031 1032 1033 1034 1035
  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 已提交
1036 1037
  tmq_list_destroy(lst);
  return rsp;
X
Xiaoyu Wang 已提交
1038 1039
}

1040 1041 1042 1043 1044 1045
static void freeClientVgImpl(void* param) {
  SMqClientTopic* pTopic = param;
  taosMemoryFreeClear(pTopic->schema.pSchema);
  taosArrayDestroy(pTopic->vgs);
}

1046
void tmqFreeImpl(void* handle) {
1047 1048
  tmq_t*  tmq = (tmq_t*)handle;
  int64_t id = tmq->consumerId;
L
Liu Jicong 已提交
1049

1050
  // TODO stop timer
L
Liu Jicong 已提交
1051 1052 1053 1054
  if (tmq->mqueue) {
    tmqClearUnhandleMsg(tmq);
    taosCloseQueue(tmq->mqueue);
  }
L
Liu Jicong 已提交
1055

H
Haojun Liao 已提交
1056 1057 1058 1059 1060
  if (tmq->delayedTask) {
    taosCloseQueue(tmq->delayedTask);
  }

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

1063
  taosArrayDestroyEx(tmq->clientTopics, freeClientVgImpl);
1064 1065
  taos_close_internal(tmq->pTscObj);
  taosMemoryFree(tmq);
1066 1067

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

1070 1071 1072 1073 1074 1075 1076 1077 1078
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);
1079
  if (tmqMgmt.rsetId < 0) {
1080 1081 1082 1083
    tmqInitRes = terrno;
  }
}

L
Liu Jicong 已提交
1084
tmq_t* tmq_consumer_new(tmq_conf_t* conf, char* errstr, int32_t errstrLen) {
1085 1086 1087 1088
  taosThreadOnce(&tmqInit, tmqMgmtInit);
  if (tmqInitRes != 0) {
    terrno = tmqInitRes;
    return NULL;
L
Liu Jicong 已提交
1089 1090
  }

L
Liu Jicong 已提交
1091 1092
  tmq_t* pTmq = taosMemoryCalloc(1, sizeof(tmq_t));
  if (pTmq == NULL) {
L
Liu Jicong 已提交
1093
    terrno = TSDB_CODE_OUT_OF_MEMORY;
1094
    tscError("failed to create consumer, groupId:%s, code:%s", conf->groupId, terrstr());
L
Liu Jicong 已提交
1095 1096
    return NULL;
  }
L
Liu Jicong 已提交
1097

L
Liu Jicong 已提交
1098 1099 1100
  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 已提交
1101 1102 1103
  pTmq->clientTopics = taosArrayInit(0, sizeof(SMqClientTopic));
  pTmq->mqueue = taosOpenQueue();
  pTmq->delayedTask = taosOpenQueue();
H
Haojun Liao 已提交
1104
  pTmq->qall = taosAllocateQall();
L
Liu Jicong 已提交
1105

X
Xiaoyu Wang 已提交
1106 1107
  if (pTmq->clientTopics == NULL || pTmq->mqueue == NULL || pTmq->qall == NULL || pTmq->delayedTask == NULL ||
      conf->groupId[0] == 0) {
L
Liu Jicong 已提交
1108
    terrno = TSDB_CODE_OUT_OF_MEMORY;
X
Xiaoyu Wang 已提交
1109
    tscError("consumer:0x%" PRIx64 " setup failed since %s, groupId:%s", pTmq->consumerId, terrstr(), pTmq->groupId);
1110
    goto _failed;
L
Liu Jicong 已提交
1111
  }
L
Liu Jicong 已提交
1112

L
Liu Jicong 已提交
1113 1114
  // init status
  pTmq->status = TMQ_CONSUMER_STATUS__INIT;
L
Liu Jicong 已提交
1115 1116
  pTmq->pollCnt = 0;
  pTmq->epoch = 0;
wmmhello's avatar
wmmhello 已提交
1117
//  pTmq->needReportOffsetRows = true;
L
Liu Jicong 已提交
1118

L
Liu Jicong 已提交
1119 1120 1121
  // set conf
  strcpy(pTmq->clientId, conf->clientId);
  strcpy(pTmq->groupId, conf->groupId);
1122
  pTmq->withTbName = conf->withTbName;
L
Liu Jicong 已提交
1123
  pTmq->useSnapshot = conf->snapEnable;
L
Liu Jicong 已提交
1124
  pTmq->autoCommit = conf->autoCommit;
L
Liu Jicong 已提交
1125
  pTmq->autoCommitInterval = conf->autoCommitInterval;
L
Liu Jicong 已提交
1126 1127
  pTmq->commitCb = conf->commitCb;
  pTmq->commitCbUserParam = conf->commitCbUserParam;
L
Liu Jicong 已提交
1128
  pTmq->resetOffsetCfg = conf->resetOffset;
wmmhello's avatar
wmmhello 已提交
1129
  taosInitRWLatch(&pTmq->lock);
L
Liu Jicong 已提交
1130

1131 1132
  pTmq->hbBgEnable = conf->hbBgEnable;

L
Liu Jicong 已提交
1133
  // assign consumerId
L
Liu Jicong 已提交
1134
  pTmq->consumerId = tGenIdPI64();
X
Xiaoyu Wang 已提交
1135

L
Liu Jicong 已提交
1136 1137
  // init semaphore
  if (tsem_init(&pTmq->rspSem, 0, 0) != 0) {
1138
    tscError("consumer:0x %" PRIx64 " setup failed since %s, consumer group %s", pTmq->consumerId, terrstr(),
S
Shengliang Guan 已提交
1139
             pTmq->groupId);
1140
    goto _failed;
L
Liu Jicong 已提交
1141
  }
L
Liu Jicong 已提交
1142

L
Liu Jicong 已提交
1143 1144 1145
  // init connection
  pTmq->pTscObj = taos_connect_internal(conf->ip, user, pass, NULL, NULL, conf->port, CONN_TYPE__TMQ);
  if (pTmq->pTscObj == NULL) {
1146
    tscError("consumer:0x%" PRIx64 " setup failed since %s, groupId:%s", pTmq->consumerId, terrstr(), pTmq->groupId);
L
Liu Jicong 已提交
1147
    tsem_destroy(&pTmq->rspSem);
1148
    goto _failed;
L
Liu Jicong 已提交
1149
  }
L
Liu Jicong 已提交
1150

1151 1152
  pTmq->refId = taosAddRef(tmqMgmt.rsetId, pTmq);
  if (pTmq->refId < 0) {
1153
    goto _failed;
1154 1155
  }

1156
  if (pTmq->hbBgEnable) {
L
Liu Jicong 已提交
1157 1158
    int64_t* pRefId = taosMemoryMalloc(sizeof(int64_t));
    *pRefId = pTmq->refId;
1159
    pTmq->hbLiveTimer = taosTmrStart(tmqSendHbReq, 1000, pRefId, tmqMgmt.timer);
1160 1161
  }

1162
  char         buf[TSDB_OFFSET_LEN] = {0};
1163 1164
  STqOffsetVal offset = {.type = pTmq->resetOffsetCfg};
  tFormatOffset(buf, tListLen(buf), &offset);
X
Xiaoyu Wang 已提交
1165 1166 1167 1168
  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 已提交
1169

1170
  return pTmq;
1171

1172 1173
_failed:
  tmqFreeImpl(pTmq);
L
Liu Jicong 已提交
1174
  return NULL;
1175 1176
}

L
Liu Jicong 已提交
1177
int32_t tmq_subscribe(tmq_t* tmq, const tmq_list_t* topic_list) {
1178
  const int32_t   MAX_RETRY_COUNT = 120 * 2;  // let's wait for 2 mins at most
L
Liu Jicong 已提交
1179 1180 1181
  const SArray*   container = &topic_list->container;
  int32_t         sz = taosArrayGetSize(container);
  void*           buf = NULL;
L
Liu Jicong 已提交
1182
  SMsgSendInfo*   sendInfo = NULL;
L
Liu Jicong 已提交
1183
  SCMSubscribeReq req = {0};
1184
  int32_t         code = 0;
1185

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

1188
  req.consumerId = tmq->consumerId;
L
Liu Jicong 已提交
1189
  tstrncpy(req.clientId, tmq->clientId, 256);
L
Liu Jicong 已提交
1190
  tstrncpy(req.cgroup, tmq->groupId, TSDB_CGROUP_LEN);
1191 1192
  req.topicNames = taosArrayInit(sz, sizeof(void*));

1193 1194 1195 1196
  if (req.topicNames == NULL) {
    code = TSDB_CODE_OUT_OF_MEMORY;
    goto FAIL;
  }
1197

1198 1199 1200 1201 1202
  req.withTbName = tmq->withTbName;
  req.autoCommit = tmq->autoCommit;
  req.autoCommitInterval = tmq->autoCommitInterval;
  req.resetOffsetCfg = tmq->resetOffsetCfg;

L
Liu Jicong 已提交
1203 1204
  for (int32_t i = 0; i < sz; i++) {
    char* topic = taosArrayGetP(container, i);
1205 1206

    SName name = {0};
L
Liu Jicong 已提交
1207 1208 1209 1210
    tNameSetDbName(&name, tmq->pTscObj->acctId, topic, strlen(topic));
    char* topicFName = taosMemoryCalloc(1, TSDB_TOPIC_FNAME_LEN);
    if (topicFName == NULL) {
      goto FAIL;
1211 1212
    }

1213
    tNameExtractFullName(&name, topicFName);
1214
    tscInfo("consumer:0x%" PRIx64 " subscribe topic:%s", tmq->consumerId, topicFName);
L
Liu Jicong 已提交
1215 1216

    taosArrayPush(req.topicNames, &topicFName);
1217 1218
  }

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

L
Liu Jicong 已提交
1221
  buf = taosMemoryMalloc(tlen);
1222 1223 1224 1225
  if (buf == NULL) {
    code = TSDB_CODE_OUT_OF_MEMORY;
    goto FAIL;
  }
L
Liu Jicong 已提交
1226

1227 1228 1229
  void* abuf = buf;
  tSerializeSCMSubscribeReq(&abuf, &req);

L
Liu Jicong 已提交
1230
  sendInfo = taosMemoryCalloc(1, sizeof(SMsgSendInfo));
1231 1232 1233 1234
  if (sendInfo == NULL) {
    code = TSDB_CODE_OUT_OF_MEMORY;
    goto FAIL;
  }
1235

H
Haojun Liao 已提交
1236
  SMqSubscribeCbParam param = { .rspErr = 0};
1237
  if (tsem_init(&param.rspSem, 0, 0) != 0) {
wmmhello's avatar
wmmhello 已提交
1238
    code = TSDB_CODE_TSC_INTERNAL_ERROR;
1239 1240
    goto FAIL;
  }
L
Liu Jicong 已提交
1241

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

L
Liu Jicong 已提交
1244 1245
  sendInfo->requestId = generateRequestId();
  sendInfo->requestObjRefId = 0;
L
Liu Jicong 已提交
1246 1247
  sendInfo->param = &param;
  sendInfo->fp = tmqSubscribeCb;
L
Liu Jicong 已提交
1248
  sendInfo->msgType = TDMT_MND_TMQ_SUBSCRIBE;
L
Liu Jicong 已提交
1249

1250 1251 1252 1253 1254
  SEpSet epSet = getEpSet_s(&tmq->pTscObj->pAppInfo->mgmtEp);

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

L
Liu Jicong 已提交
1255 1256
  // avoid double free if msg is sent
  buf = NULL;
L
Liu Jicong 已提交
1257
  sendInfo = NULL;
L
Liu Jicong 已提交
1258

L
Liu Jicong 已提交
1259 1260
  tsem_wait(&param.rspSem);
  tsem_destroy(&param.rspSem);
1261

1262 1263 1264 1265
  if (param.rspErr != 0) {
    code = param.rspErr;
    goto FAIL;
  }
L
Liu Jicong 已提交
1266

L
Liu Jicong 已提交
1267
  int32_t retryCnt = 0;
1268
  while (TSDB_CODE_MND_CONSUMER_NOT_READY == doAskEp(tmq)) {
1269
    if (retryCnt++ > MAX_RETRY_COUNT) {
wmmhello's avatar
wmmhello 已提交
1270
      tscError("consumer:0x%" PRIx64 ", mnd not ready for subscribe, retry:%d in 500ms", tmq->consumerId, retryCnt);
1271
      code = TSDB_CODE_MND_CONSUMER_NOT_READY;
L
Liu Jicong 已提交
1272 1273
      goto FAIL;
    }
1274

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

1279 1280
  // init ep timer
  if (tmq->epTimer == NULL) {
1281 1282 1283
    int64_t* pRefId1 = taosMemoryMalloc(sizeof(int64_t));
    *pRefId1 = tmq->refId;
    tmq->epTimer = taosTmrStart(tmqAssignAskEpTask, 1000, pRefId1, tmqMgmt.timer);
1284
  }
L
Liu Jicong 已提交
1285 1286

  // init auto commit timer
1287
  if (tmq->autoCommit && tmq->commitTimer == NULL) {
1288 1289 1290
    int64_t* pRefId2 = taosMemoryMalloc(sizeof(int64_t));
    *pRefId2 = tmq->refId;
    tmq->commitTimer = taosTmrStart(tmqAssignDelayedCommitTask, tmq->autoCommitInterval, pRefId2, tmqMgmt.timer);
L
Liu Jicong 已提交
1291 1292
  }

L
Liu Jicong 已提交
1293
FAIL:
L
Liu Jicong 已提交
1294
  taosArrayDestroyP(req.topicNames, taosMemoryFree);
L
Liu Jicong 已提交
1295
  taosMemoryFree(buf);
L
Liu Jicong 已提交
1296
  taosMemoryFree(sendInfo);
L
Liu Jicong 已提交
1297

L
Liu Jicong 已提交
1298
  return code;
1299 1300
}

L
Liu Jicong 已提交
1301
void tmq_conf_set_auto_commit_cb(tmq_conf_t* conf, tmq_commit_cb* cb, void* param) {
1302
  conf->commitCb = cb;
L
Liu Jicong 已提交
1303
  conf->commitCbUserParam = param;
L
Liu Jicong 已提交
1304
}
1305

wmmhello's avatar
wmmhello 已提交
1306 1307 1308 1309 1310 1311 1312 1313 1314 1315 1316 1317 1318 1319 1320 1321 1322 1323 1324 1325 1326 1327 1328 1329 1330 1331 1332 1333
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 已提交
1334
int32_t tmqPollCb(void* param, SDataBuf* pMsg, int32_t code) {
X
Xiaoyu Wang 已提交
1335
  SMqPollCbParam* pParam = (SMqPollCbParam*)param;
1336 1337

  int64_t         refId = pParam->refId;
wmmhello's avatar
wmmhello 已提交
1338 1339
//  SMqClientVg*    pVg = pParam->pVg;
//  SMqClientTopic* pTopic = pParam->pTopic;
1340

1341
  tmq_t* tmq = taosAcquireRef(tmqMgmt.rsetId, refId);
1342 1343 1344
  if (tmq == NULL) {
    taosMemoryFree(pParam);
    taosMemoryFree(pMsg->pData);
dengyihao's avatar
dengyihao 已提交
1345
    taosMemoryFree(pMsg->pEpSet);
1346 1347 1348 1349
    terrno = TSDB_CODE_TMQ_CONSUMER_CLOSED;
    return -1;
  }

H
Haojun Liao 已提交
1350 1351 1352 1353
  int32_t  epoch = pParam->epoch;
  int32_t  vgId = pParam->vgId;
  uint64_t requestId = pParam->requestId;

L
Liu Jicong 已提交
1354
  if (code != 0) {
L
Liu Jicong 已提交
1355
    if (pMsg->pData) taosMemoryFree(pMsg->pData);
dengyihao's avatar
dengyihao 已提交
1356 1357
    if (pMsg->pEpSet) taosMemoryFree(pMsg->pEpSet);

H
Haojun Liao 已提交
1358
    // in case of consumer mismatch, wait for 500ms and retry
L
Liu Jicong 已提交
1359
    if (code == TSDB_CODE_TMQ_CONSUMER_MISMATCH) {
1360
//      taosMsleep(500);
L
Liu Jicong 已提交
1361
      atomic_store_8(&tmq->status, TMQ_CONSUMER_STATUS__RECOVER);
X
Xiaoyu Wang 已提交
1362 1363
      tscDebug("consumer:0x%" PRIx64 " wait for the re-balance, wait for 500ms and set status to be RECOVER",
               tmq->consumerId);
H
Haojun Liao 已提交
1364
    } else if (code == TSDB_CODE_TQ_NO_COMMITTED_OFFSET) {
S
Shengliang Guan 已提交
1365
      SMqPollRspWrapper* pRspWrapper = taosAllocateQitem(sizeof(SMqPollRspWrapper), DEF_QITEM, 0);
L
Liu Jicong 已提交
1366
      if (pRspWrapper == NULL) {
H
Haojun Liao 已提交
1367 1368
        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 已提交
1369 1370
        goto CREATE_MSG_FAIL;
      }
H
Haojun Liao 已提交
1371

L
Liu Jicong 已提交
1372 1373
      pRspWrapper->tmqRspType = TMQ_MSG_TYPE__END_RSP;
      taosWriteQitem(tmq->mqueue, pRspWrapper);
1374 1375
//    } else if (code == TSDB_CODE_WAL_LOG_NOT_EXIST) {  // poll data while insert
//      taosMsleep(5);
wmmhello's avatar
wmmhello 已提交
1376 1377 1378
    } 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 已提交
1379
    }
H
Haojun Liao 已提交
1380

L
fix txn  
Liu Jicong 已提交
1381
    goto CREATE_MSG_FAIL;
L
Liu Jicong 已提交
1382 1383
  }

X
Xiaoyu Wang 已提交
1384
  int32_t msgEpoch = ((SMqRspHead*)pMsg->pData)->epoch;
1385 1386
  int32_t clientEpoch = atomic_load_32(&tmq->epoch);
  if (msgEpoch < clientEpoch) {
L
Liu Jicong 已提交
1387
    // do not write into queue since updating epoch reset
X
Xiaoyu Wang 已提交
1388 1389
    tscWarn("consumer:0x%" PRIx64
            " msg discard from vgId:%d since from earlier epoch, rsp epoch %d, current epoch %d, reqId:0x%" PRIx64,
1390
            tmq->consumerId, vgId, msgEpoch, clientEpoch, requestId);
H
Haojun Liao 已提交
1391

1392
    tsem_post(&tmq->rspSem);
1393 1394
    taosReleaseRef(tmqMgmt.rsetId, refId);

L
Liu Jicong 已提交
1395
    taosMemoryFree(pMsg->pData);
dengyihao's avatar
dengyihao 已提交
1396
    taosMemoryFree(pMsg->pEpSet);
wmmhello's avatar
wmmhello 已提交
1397 1398
    taosMemoryFree(pParam);

X
Xiaoyu Wang 已提交
1399 1400 1401
    return 0;
  }

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

L
Liu Jicong 已提交
1407 1408 1409
  // handle meta rsp
  int8_t rspType = ((SMqRspHead*)pMsg->pData)->mqMsgType;

S
Shengliang Guan 已提交
1410
  SMqPollRspWrapper* pRspWrapper = taosAllocateQitem(sizeof(SMqPollRspWrapper), DEF_QITEM, 0);
L
Liu Jicong 已提交
1411
  if (pRspWrapper == NULL) {
L
Liu Jicong 已提交
1412
    taosMemoryFree(pMsg->pData);
dengyihao's avatar
dengyihao 已提交
1413
    taosMemoryFree(pMsg->pEpSet);
X
Xiaoyu Wang 已提交
1414 1415
    tscWarn("consumer:0x%" PRIx64 " msg discard from vgId:%d, epoch %d since out of memory", tmq->consumerId, vgId,
            epoch);
L
fix txn  
Liu Jicong 已提交
1416
    goto CREATE_MSG_FAIL;
L
Liu Jicong 已提交
1417
  }
L
Liu Jicong 已提交
1418

L
Liu Jicong 已提交
1419
  pRspWrapper->tmqRspType = rspType;
wmmhello's avatar
wmmhello 已提交
1420 1421
//  pRspWrapper->vgHandle = pVg;
//  pRspWrapper->topicHandle = pTopic;
H
Haojun Liao 已提交
1422
  pRspWrapper->reqId = requestId;
1423
  pRspWrapper->pEpset = pMsg->pEpSet;
wmmhello's avatar
wmmhello 已提交
1424 1425
  pRspWrapper->vgId = vgId;
  strcpy(pRspWrapper->topicName, pParam->topicName);
L
Liu Jicong 已提交
1426

1427
  pMsg->pEpSet = NULL;
L
Liu Jicong 已提交
1428
  if (rspType == TMQ_MSG_TYPE__POLL_RSP) {
L
Liu Jicong 已提交
1429 1430
    SDecoder decoder;
    tDecoderInit(&decoder, POINTER_SHIFT(pMsg->pData, sizeof(SMqRspHead)), pMsg->len - sizeof(SMqRspHead));
1431
    tDecodeMqDataRsp(&decoder, &pRspWrapper->dataRsp);
wmmhello's avatar
wmmhello 已提交
1432
    tDecoderClear(&decoder);
L
Liu Jicong 已提交
1433
    memcpy(&pRspWrapper->dataRsp, pMsg->pData, sizeof(SMqRspHead));
L
Liu Jicong 已提交
1434

wmmhello's avatar
wmmhello 已提交
1435
    char buf[TSDB_OFFSET_LEN] = {0};
1436
    tFormatOffset(buf, TSDB_OFFSET_LEN, &pRspWrapper->dataRsp.rspOffset);
H
Haojun Liao 已提交
1437
    tscDebug("consumer:0x%" PRIx64 " recv poll rsp, vgId:%d, req ver:%" PRId64 ", rsp:%s type %d, reqId:0x%" PRIx64,
H
Haojun Liao 已提交
1438
             tmq->consumerId, vgId, pRspWrapper->dataRsp.reqOffset.version, buf, rspType, requestId);
L
Liu Jicong 已提交
1439
  } else if (rspType == TMQ_MSG_TYPE__POLL_META_RSP) {
1440 1441
    SDecoder decoder;
    tDecoderInit(&decoder, POINTER_SHIFT(pMsg->pData, sizeof(SMqRspHead)), pMsg->len - sizeof(SMqRspHead));
1442
    tDecodeMqMetaRsp(&decoder, &pRspWrapper->metaRsp);
1443
    tDecoderClear(&decoder);
L
Liu Jicong 已提交
1444
    memcpy(&pRspWrapper->metaRsp, pMsg->pData, sizeof(SMqRspHead));
L
Liu Jicong 已提交
1445 1446 1447 1448 1449 1450
  } 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 已提交
1451 1452
  } else {  // invalid rspType
    tscError("consumer:0x%" PRIx64 " invalid rsp msg received, type:%d ignored", tmq->consumerId, rspType);
L
Liu Jicong 已提交
1453
  }
L
Liu Jicong 已提交
1454

L
Liu Jicong 已提交
1455
  taosMemoryFree(pMsg->pData);
H
Haojun Liao 已提交
1456
  taosWriteQitem(tmq->mqueue, pRspWrapper);
L
Liu Jicong 已提交
1457

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

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

L
Liu Jicong 已提交
1466
  return 0;
H
Haojun Liao 已提交
1467

L
fix txn  
Liu Jicong 已提交
1468
CREATE_MSG_FAIL:
L
Liu Jicong 已提交
1469
  if (epoch == tmq->epoch) {
wmmhello's avatar
wmmhello 已提交
1470 1471 1472 1473 1474
    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 已提交
1475
    taosWUnLockLatch(&tmq->lock);
L
Liu Jicong 已提交
1476
  }
H
Haojun Liao 已提交
1477

1478
  tsem_post(&tmq->rspSem);
1479
  taosReleaseRef(tmqMgmt.rsetId, refId);
wmmhello's avatar
wmmhello 已提交
1480
  taosMemoryFree(pParam);
1481

L
Liu Jicong 已提交
1482
  return -1;
1483 1484
}

H
Haojun Liao 已提交
1485
typedef struct SVgroupSaveInfo {
wmmhello's avatar
wmmhello 已提交
1486 1487
  STqOffsetVal currentOffset;
  STqOffsetVal commitOffset;
1488
  STqOffsetVal seekOffset;
H
Haojun Liao 已提交
1489 1490 1491
  int64_t      numOfRows;
} SVgroupSaveInfo;

H
Haojun Liao 已提交
1492 1493 1494 1495 1496 1497
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 已提交
1498
  char    vgKey[TSDB_TOPIC_FNAME_LEN + 22];
H
Haojun Liao 已提交
1499 1500 1501 1502 1503
  int32_t vgNumGet = taosArrayGetSize(pTopicEp->vgs);

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

1504
  tscInfo("consumer:0x%" PRIx64 ", update topic:%s, new numOfVgs:%d", tmq->consumerId, pTopic->topicName, vgNumGet);
H
Haojun Liao 已提交
1505 1506 1507 1508
  pTopic->vgs = taosArrayInit(vgNumGet, sizeof(SMqClientVg));

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

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

wmmhello's avatar
wmmhello 已提交
1513 1514
    STqOffsetVal offsetNew = {0};
    offsetNew.type = tmq->resetOffsetCfg;
H
Haojun Liao 已提交
1515 1516 1517 1518 1519

    SMqClientVg clientVg = {
        .pollCnt = 0,
        .vgId = pVgEp->vgId,
        .epSet = pVgEp->epSet,
wmmhello's avatar
wmmhello 已提交
1520
        .vgStatus = TMQ_VG_STATUS__IDLE,
H
Haojun Liao 已提交
1521
        .vgSkipCnt = 0,
H
Haojun Liao 已提交
1522
        .emptyBlockReceiveTs = 0,
wmmhello's avatar
wmmhello 已提交
1523
        .numOfRows = pInfo ? pInfo->numOfRows : 0,
H
Haojun Liao 已提交
1524 1525
    };

wmmhello's avatar
wmmhello 已提交
1526 1527
    clientVg.offsetInfo.currentOffset = pInfo ? pInfo->currentOffset : offsetNew;
    clientVg.offsetInfo.committedOffset = pInfo ? pInfo->commitOffset : offsetNew;
1528
    clientVg.offsetInfo.seekOffset = pInfo ? pInfo->seekOffset : offsetNew;
H
Haojun Liao 已提交
1529 1530
    clientVg.offsetInfo.walVerBegin = -1;
    clientVg.offsetInfo.walVerEnd = -1;
1531
    clientVg.seekUpdated = false;
1532
//    clientVg.receivedInfoFromVnode = false;
1533

H
Haojun Liao 已提交
1534 1535 1536 1537 1538 1539 1540 1541 1542 1543 1544 1545 1546
    taosArrayPush(pTopic->vgs, &clientVg);
  }
}

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

  taosArrayDestroy(pTopic->vgs);
}

1547
static bool doUpdateLocalEp(tmq_t* tmq, int32_t epoch, const SMqAskEpRsp* pRsp) {
1548 1549 1550
  bool set = false;

  int32_t topicNumGet = taosArrayGetSize(pRsp->topics);
wmmhello's avatar
wmmhello 已提交
1551 1552 1553
  if (epoch <= tmq->epoch) {
    return false;
  }
1554 1555 1556 1557 1558 1559

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

H
Haojun Liao 已提交
1560 1561
  SHashObj* pVgOffsetHashMap = taosHashInit(64, MurmurHash3_32, false, HASH_NO_LOCK);
  if (pVgOffsetHashMap == NULL) {
1562 1563 1564
    taosArrayDestroy(newTopics);
    return false;
  }
1565

wmmhello's avatar
wmmhello 已提交
1566 1567 1568 1569 1570 1571
  taosWLockLatch(&tmq->lock);
  int32_t topicNumCur = taosArrayGetSize(tmq->clientTopics);

  char vgKey[TSDB_TOPIC_FNAME_LEN + 22];
  tscInfo("consumer:0x%" PRIx64 " update ep epoch from %d to epoch %d, incoming topics:%d, existed topics:%d",
          tmq->consumerId, tmq->epoch, epoch, topicNumGet, topicNumCur);
H
Haojun Liao 已提交
1572
  // todo extract method
1573 1574 1575 1576 1577
  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);
1578
      tscInfo("consumer:0x%" PRIx64 ", current vg num: %d", tmq->consumerId, vgNumCur);
1579 1580
      for (int32_t j = 0; j < vgNumCur; j++) {
        SMqClientVg* pVgCur = taosArrayGet(pTopicCur->vgs, j);
H
Haojun Liao 已提交
1581 1582
        makeTopicVgroupKey(vgKey, pTopicCur->topicName, pVgCur->vgId);

wmmhello's avatar
wmmhello 已提交
1583
        char buf[TSDB_OFFSET_LEN] = {0};
1584
        tFormatOffset(buf, TSDB_OFFSET_LEN, &pVgCur->offsetInfo.currentOffset);
1585
        tscInfo("consumer:0x%" PRIx64 ", epoch:%d vgId:%d vgKey:%s, offset:%s", tmq->consumerId, epoch, pVgCur->vgId,
X
Xiaoyu Wang 已提交
1586
                 vgKey, buf);
H
Haojun Liao 已提交
1587

1588
        SVgroupSaveInfo info = {.currentOffset = pVgCur->offsetInfo.currentOffset, .seekOffset = pVgCur->offsetInfo.seekOffset, .commitOffset = pVgCur->offsetInfo.committedOffset, .numOfRows = pVgCur->numOfRows};
H
Haojun Liao 已提交
1589
        taosHashPut(pVgOffsetHashMap, vgKey, strlen(vgKey), &info, sizeof(SVgroupSaveInfo));
1590 1591 1592 1593 1594 1595 1596
      }
    }
  }

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

H
Haojun Liao 已提交
1601 1602
  taosHashCleanup(pVgOffsetHashMap);

1603
  // destroy current buffered existed topics info
1604
  if (tmq->clientTopics) {
H
Haojun Liao 已提交
1605
    taosArrayDestroyEx(tmq->clientTopics, freeClientVgInfo);
X
Xiaoyu Wang 已提交
1606
  }
H
Haojun Liao 已提交
1607
  tmq->clientTopics = newTopics;
wmmhello's avatar
wmmhello 已提交
1608
  taosWUnLockLatch(&tmq->lock);
1609

X
Xiaoyu Wang 已提交
1610
  int8_t flag = (topicNumGet == 0) ? TMQ_CONSUMER_STATUS__NO_TOPIC : TMQ_CONSUMER_STATUS__READY;
H
Haojun Liao 已提交
1611
  atomic_store_8(&tmq->status, flag);
X
Xiaoyu Wang 已提交
1612
  atomic_store_32(&tmq->epoch, epoch);
H
Haojun Liao 已提交
1613

1614
  tscInfo("consumer:0x%" PRIx64 " update topic info completed", tmq->consumerId);
X
Xiaoyu Wang 已提交
1615 1616 1617
  return set;
}

1618
int32_t askEpCallbackFn(void* param, SDataBuf* pMsg, int32_t code) {
1619
  SMqAskEpCbParam* pParam = (SMqAskEpCbParam*)param;
1620 1621 1622
  tmq_t*           tmq = taosAcquireRef(tmqMgmt.rsetId, pParam->refId);

  if (tmq == NULL) {
1623
    terrno = TSDB_CODE_TMQ_CONSUMER_CLOSED;
wmmhello's avatar
wmmhello 已提交
1624
//    pParam->pUserFn(tmq, terrno, NULL, pParam->pParam);
1625

1626
    taosMemoryFree(pMsg->pData);
dengyihao's avatar
dengyihao 已提交
1627
    taosMemoryFree(pMsg->pEpSet);
1628 1629
    taosMemoryFree(pParam);
    return terrno;
1630 1631
  }

H
Haojun Liao 已提交
1632
  if (code != TSDB_CODE_SUCCESS) {
1633 1634 1635 1636 1637 1638 1639 1640 1641
    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;
1642
  }
L
Liu Jicong 已提交
1643

L
Liu Jicong 已提交
1644
  // tmq's epoch is monotonically increase,
L
Liu Jicong 已提交
1645
  // so it's safe to discard any old epoch msg.
L
Liu Jicong 已提交
1646
  // Epoch will only increase when received newer epoch ep msg
L
Liu Jicong 已提交
1647 1648 1649
  SMqRspHead* head = pMsg->pData;
  int32_t     epoch = atomic_load_32(&tmq->epoch);
  if (head->epoch <= epoch) {
1650
    tscInfo("consumer:0x%" PRIx64 ", recv ep, msg epoch %d, current epoch %d, no need to update local ep",
1651
             tmq->consumerId, head->epoch, epoch);
1652

1653 1654 1655 1656 1657 1658 1659 1660
    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 已提交
1661
  } else {
1662
    tscInfo("consumer:0x%" PRIx64 ", recv ep, msg epoch %d, current epoch %d, update local ep", tmq->consumerId,
1663
             head->epoch, epoch);
1664
  }
L
Liu Jicong 已提交
1665

1666
  pParam->pUserFn(tmq, code, pMsg, pParam->pParam);
1667 1668
  taosReleaseRef(tmqMgmt.rsetId, pParam->refId);

dengyihao's avatar
dengyihao 已提交
1669
  taosMemoryFree(pMsg->pEpSet);
L
Liu Jicong 已提交
1670
  taosMemoryFree(pMsg->pData);
1671
  taosMemoryFree(pParam);
L
Liu Jicong 已提交
1672
  return code;
1673 1674
}

L
Liu Jicong 已提交
1675
void tmqBuildConsumeReqImpl(SMqPollReq* pReq, tmq_t* tmq, int64_t timeout, SMqClientTopic* pTopic, SMqClientVg* pVg) {
L
Liu Jicong 已提交
1676 1677 1678 1679
  int32_t groupLen = strlen(tmq->groupId);
  memcpy(pReq->subKey, tmq->groupId, groupLen);
  pReq->subKey[groupLen] = TMQ_SEPARATOR;
  strcpy(pReq->subKey + groupLen + 1, pTopic->topicName);
1680

1681
  pReq->withTbName = tmq->withTbName;
L
Liu Jicong 已提交
1682
  pReq->consumerId = tmq->consumerId;
1683
  pReq->timeout = timeout;
X
Xiaoyu Wang 已提交
1684
  pReq->epoch = tmq->epoch;
H
Haojun Liao 已提交
1685
  pReq->reqOffset = pVg->offsetInfo.currentOffset;
D
dapan1121 已提交
1686
  pReq->head.vgId = pVg->vgId;
1687 1688
  pReq->useSnapshot = tmq->useSnapshot;
  pReq->reqId = generateRequestId();
1689 1690
}

L
Liu Jicong 已提交
1691 1692
SMqMetaRspObj* tmqBuildMetaRspFromWrapper(SMqPollRspWrapper* pWrapper) {
  SMqMetaRspObj* pRspObj = taosMemoryCalloc(1, sizeof(SMqMetaRspObj));
L
Liu Jicong 已提交
1693
  pRspObj->resType = RES_TYPE__TMQ_META;
L
Liu Jicong 已提交
1694 1695 1696 1697 1698 1699 1700 1701
  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;
}

1702
SMqRspObj* tmqBuildRspFromWrapper(SMqPollRspWrapper* pWrapper, SMqClientVg* pVg, int64_t* numOfRows) {
L
Liu Jicong 已提交
1703 1704
  SMqRspObj* pRspObj = taosMemoryCalloc(1, sizeof(SMqRspObj));
  pRspObj->resType = RES_TYPE__TMQ;
1705

1706
  (*numOfRows) = 0;
L
Liu Jicong 已提交
1707 1708
  tstrncpy(pRspObj->topic, pWrapper->topicHandle->topicName, TSDB_TOPIC_FNAME_LEN);
  tstrncpy(pRspObj->db, pWrapper->topicHandle->db, TSDB_DB_FNAME_LEN);
1709

L
Liu Jicong 已提交
1710
  pRspObj->vgId = pWrapper->vgHandle->vgId;
L
Liu Jicong 已提交
1711
  pRspObj->resIter = -1;
L
Liu Jicong 已提交
1712
  memcpy(&pRspObj->rsp, &pWrapper->dataRsp, sizeof(SMqDataRsp));
L
Liu Jicong 已提交
1713

L
Liu Jicong 已提交
1714 1715
  pRspObj->resInfo.totalRows = 0;
  pRspObj->resInfo.precision = TSDB_TIME_PRECISION_MILLI;
1716

L
Liu Jicong 已提交
1717
  if (!pWrapper->dataRsp.withSchema) {
L
Liu Jicong 已提交
1718 1719
    setResSchemaInfo(&pRspObj->resInfo, pWrapper->topicHandle->schema.pSchema, pWrapper->topicHandle->schema.nCols);
  }
L
Liu Jicong 已提交
1720

1721
  // extract the rows in this data packet
X
Xiaoyu Wang 已提交
1722
  for (int32_t i = 0; i < pRspObj->rsp.blockNum; ++i) {
1723
    SRetrieveTableRsp* pRetrieve = (SRetrieveTableRsp*)taosArrayGetP(pRspObj->rsp.blockData, i);
X
Xiaoyu Wang 已提交
1724
    int64_t            rows = htobe64(pRetrieve->numOfRows);
1725
    pVg->numOfRows += rows;
1726
    (*numOfRows) += rows;
1727 1728
  }

L
Liu Jicong 已提交
1729
  return pRspObj;
X
Xiaoyu Wang 已提交
1730 1731
}

1732
SMqTaosxRspObj* tmqBuildTaosxRspFromWrapper(SMqPollRspWrapper* pWrapper, SMqClientVg* pVg, int64_t* numOfRows) {
L
Liu Jicong 已提交
1733
  SMqTaosxRspObj* pRspObj = taosMemoryCalloc(1, sizeof(SMqTaosxRspObj));
1734
  pRspObj->resType = RES_TYPE__TMQ_METADATA;
L
Liu Jicong 已提交
1735 1736 1737 1738
  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;
1739
  memcpy(&pRspObj->rsp, &pWrapper->taosxRsp, sizeof(STaosxRsp));
L
Liu Jicong 已提交
1740 1741 1742

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

1747 1748 1749 1750 1751 1752 1753
  // 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 已提交
1754 1755 1756
  return pRspObj;
}

1757 1758 1759 1760 1761 1762 1763 1764 1765 1766 1767 1768 1769 1770 1771 1772 1773 1774 1775 1776 1777 1778 1779 1780 1781 1782 1783 1784 1785 1786 1787 1788 1789
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 已提交
1790 1791 1792
//  pParam->pVg = pVg;  // pVg may be released,fix it
//  pParam->pTopic = pTopic;
  strcpy(pParam->topicName, pTopic->topicName);
1793
  pParam->vgId = pVg->vgId;
H
Haojun Liao 已提交
1794
  pParam->requestId = req.reqId;
1795 1796 1797 1798 1799 1800 1801 1802

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

H
Haojun Liao 已提交
1803
  sendInfo->msgInfo = (SDataBuf){.pData = msg, .len = msgSize, .handle = NULL};
1804 1805 1806 1807 1808 1809 1810
  sendInfo->requestId = req.reqId;
  sendInfo->requestObjRefId = 0;
  sendInfo->param = pParam;
  sendInfo->fp = tmqPollCb;
  sendInfo->msgType = TDMT_VND_TMQ_CONSUME;

  int64_t transporterId = 0;
wmmhello's avatar
wmmhello 已提交
1811
  char    offsetFormatBuf[TSDB_OFFSET_LEN] = {0};
H
Haojun Liao 已提交
1812
  tFormatOffset(offsetFormatBuf, tListLen(offsetFormatBuf), &pVg->offsetInfo.currentOffset);
1813

X
Xiaoyu Wang 已提交
1814 1815
  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);
1816 1817 1818
  asyncSendMsgToServer(pTmq->pTscObj->pAppInfo->pTransporter, &pVg->epSet, &transporterId, sendInfo);

  pVg->pollCnt++;
1819
  pVg->seekUpdated = false;   // reset this flag.
1820 1821 1822 1823 1824
  pTmq->pollCnt++;

  return TSDB_CODE_SUCCESS;
}

1825
// broadcast the poll request to all related vnodes
H
Haojun Liao 已提交
1826
static int32_t tmqPollImpl(tmq_t* tmq, int64_t timeout) {
1827 1828 1829
  if(atomic_load_8(&tmq->status) == TMQ_CONSUMER_STATUS__RECOVER){
    return 0;
  }
wmmhello's avatar
wmmhello 已提交
1830 1831 1832
  int32_t code = 0;

  taosWLockLatch(&tmq->lock);
1833
  int32_t numOfTopics = taosArrayGetSize(tmq->clientTopics);
X
Xiaoyu Wang 已提交
1834
  tscDebug("consumer:0x%" PRIx64 " start to poll data, numOfTopics:%d", tmq->consumerId, numOfTopics);
1835 1836

  for (int i = 0; i < numOfTopics; i++) {
X
Xiaoyu Wang 已提交
1837
    SMqClientTopic* pTopic = taosArrayGet(tmq->clientTopics, i);
X
Xiaoyu Wang 已提交
1838
    int32_t         numOfVg = taosArrayGetSize(pTopic->vgs);
1839 1840

    for (int j = 0; j < numOfVg; j++) {
X
Xiaoyu Wang 已提交
1841
      SMqClientVg* pVg = taosArrayGet(pTopic->vgs, j);
wmmhello's avatar
wmmhello 已提交
1842
      if (taosGetTimestampMs() - pVg->emptyBlockReceiveTs < EMPTY_BLOCK_POLL_IDLE_DURATION) {  // less than 10ms
1843
        tscTrace("consumer:0x%" PRIx64 " epoch %d, vgId:%d idle for 10ms before start next poll", tmq->consumerId,
X
Xiaoyu Wang 已提交
1844
                 tmq->epoch, pVg->vgId);
H
Haojun Liao 已提交
1845 1846 1847
        continue;
      }

1848
      int32_t vgStatus = atomic_val_compare_exchange_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE, TMQ_VG_STATUS__WAIT);
1849
      if (vgStatus == TMQ_VG_STATUS__WAIT) {
L
Liu Jicong 已提交
1850
        int32_t vgSkipCnt = atomic_add_fetch_32(&pVg->vgSkipCnt, 1);
1851
        tscTrace("consumer:0x%" PRIx64 " epoch %d wait poll-rsp, skip vgId:%d skip cnt %d", tmq->consumerId, tmq->epoch,
X
Xiaoyu Wang 已提交
1852
                 pVg->vgId, vgSkipCnt);
X
Xiaoyu Wang 已提交
1853 1854
        continue;
      }
1855

L
Liu Jicong 已提交
1856
      atomic_store_32(&pVg->vgSkipCnt, 0);
wmmhello's avatar
wmmhello 已提交
1857
      code = doTmqPollImpl(tmq, pTopic, pVg, timeout);
1858
      if (code != TSDB_CODE_SUCCESS) {
wmmhello's avatar
wmmhello 已提交
1859
        goto end;
D
dapan1121 已提交
1860
      }
X
Xiaoyu Wang 已提交
1861 1862
    }
  }
1863

wmmhello's avatar
wmmhello 已提交
1864 1865 1866 1867
end:
  taosWUnLockLatch(&tmq->lock);
  tscDebug("consumer:0x%" PRIx64 " end to poll data, code:%d", tmq->consumerId, code);
  return code;
X
Xiaoyu Wang 已提交
1868 1869
}

H
Haojun Liao 已提交
1870
static int32_t tmqHandleNoPollRsp(tmq_t* tmq, SMqRspWrapper* rspWrapper, bool* pReset) {
L
Liu Jicong 已提交
1871
  if (rspWrapper->tmqRspType == TMQ_MSG_TYPE__EP_RSP) {
L
fix  
Liu Jicong 已提交
1872
    /*printf("ep %d %d\n", rspMsg->head.epoch, tmq->epoch);*/
L
Liu Jicong 已提交
1873 1874
    if (rspWrapper->epoch > atomic_load_32(&tmq->epoch)) {
      SMqAskEpRspWrapper* pEpRspWrapper = (SMqAskEpRspWrapper*)rspWrapper;
L
Liu Jicong 已提交
1875
      SMqAskEpRsp*        rspMsg = &pEpRspWrapper->msg;
1876
      doUpdateLocalEp(tmq, rspWrapper->epoch, rspMsg);
L
temp  
Liu Jicong 已提交
1877
      /*tmqClearUnhandleMsg(tmq);*/
L
Liu Jicong 已提交
1878
      tDeleteSMqAskEpRsp(rspMsg);
X
Xiaoyu Wang 已提交
1879 1880
      *pReset = true;
    } else {
L
Liu Jicong 已提交
1881
      tmqFreeRspWrapper(rspWrapper);
X
Xiaoyu Wang 已提交
1882 1883 1884 1885 1886 1887 1888 1889
      *pReset = false;
    }
  } else {
    return -1;
  }
  return 0;
}

1890
static void updateVgInfo(SMqClientVg* pVg, STqOffsetVal* reqOffset, STqOffsetVal* rspOffset, int64_t sver, int64_t ever, int64_t consumerId){
wmmhello's avatar
wmmhello 已提交
1891 1892
  if (!pVg->seekUpdated) {
    tscDebug("consumer:0x%" PRIx64" local offset is update, since seekupdate not set", consumerId);
1893 1894
    pVg->offsetInfo.seekOffset = *reqOffset;
    pVg->offsetInfo.currentOffset = *rspOffset;
wmmhello's avatar
wmmhello 已提交
1895 1896 1897 1898 1899 1900 1901 1902 1903 1904
  } else {
    tscDebug("consumer:0x%" PRIx64" local offset is NOT update, since seekupdate is set", consumerId);
  }

  // update the status
  atomic_store_32(&pVg->vgStatus, TMQ_VG_STATUS__IDLE);

  // update the valid wal version range
  pVg->offsetInfo.walVerBegin = sver;
  pVg->offsetInfo.walVerEnd = ever;
1905
//  pVg->receivedInfoFromVnode = true;
wmmhello's avatar
wmmhello 已提交
1906 1907
}

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

X
Xiaoyu Wang 已提交
1911
  while (1) {
1912 1913
    SMqRspWrapper* pRspWrapper = NULL;
    taosGetQitem(tmq->qall, (void**)&pRspWrapper);
1914

1915
    if (pRspWrapper == NULL) {
L
Liu Jicong 已提交
1916
      taosReadAllQitems(tmq->mqueue, tmq->qall);
1917 1918
      taosGetQitem(tmq->qall, (void**)&pRspWrapper);
      if (pRspWrapper == NULL) {
L
Liu Jicong 已提交
1919 1920
        return NULL;
      }
X
Xiaoyu Wang 已提交
1921 1922
    }

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

1925 1926
    if (pRspWrapper->tmqRspType == TMQ_MSG_TYPE__END_RSP) {
      taosFreeQitem(pRspWrapper);
L
Liu Jicong 已提交
1927
      terrno = TSDB_CODE_TQ_NO_COMMITTED_OFFSET;
H
Haojun Liao 已提交
1928
      tscError("consumer:0x%" PRIx64 " unexpected rsp from poll, code:%s", tmq->consumerId, tstrerror(terrno));
L
Liu Jicong 已提交
1929
      return NULL;
1930 1931
    } else if (pRspWrapper->tmqRspType == TMQ_MSG_TYPE__POLL_RSP) {
      SMqPollRspWrapper* pollRspWrapper = (SMqPollRspWrapper*)pRspWrapper;
H
Haojun Liao 已提交
1932

X
Xiaoyu Wang 已提交
1933
      int32_t     consumerEpoch = atomic_load_32(&tmq->epoch);
1934 1935 1936
      SMqDataRsp* pDataRsp = &pollRspWrapper->dataRsp;

      if (pDataRsp->head.epoch == consumerEpoch) {
wmmhello's avatar
wmmhello 已提交
1937
        taosWLockLatch(&tmq->lock);
wmmhello's avatar
wmmhello 已提交
1938 1939 1940 1941 1942 1943
        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);
wmmhello's avatar
wmmhello 已提交
1944
          taosWUnLockLatch(&tmq->lock);
wmmhello's avatar
wmmhello 已提交
1945 1946
          return NULL;
        }
1947 1948 1949 1950
        // update the epset
        if (pollRspWrapper->pEpset != NULL) {
          SEp* pEp = GET_ACTIVE_EP(pollRspWrapper->pEpset);
          SEp* pOld = GET_ACTIVE_EP(&(pVg->epSet));
X
Xiaoyu Wang 已提交
1951 1952
          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);
1953 1954 1955
          pVg->epSet = *pollRspWrapper->pEpset;
        }

1956
        updateVgInfo(pVg, &pDataRsp->reqOffset, &pDataRsp->rspOffset, pDataRsp->head.walsver, pDataRsp->head.walever, tmq->consumerId);
1957

wmmhello's avatar
wmmhello 已提交
1958
        char buf[TSDB_OFFSET_LEN] = {0};
1959
        tFormatOffset(buf, TSDB_OFFSET_LEN, &pDataRsp->rspOffset);
1960
        if (pDataRsp->blockNum == 0) {
X
Xiaoyu Wang 已提交
1961
          tscDebug("consumer:0x%" PRIx64 " empty block received, vgId:%d, offset:%s, vg total:%" PRId64
wmmhello's avatar
wmmhello 已提交
1962
                   ", total:%" PRId64 ", reqId:0x%" PRIx64,
X
Xiaoyu Wang 已提交
1963
                   tmq->consumerId, pVg->vgId, buf, pVg->numOfRows, tmq->totalRows, pollRspWrapper->reqId);
1964
          pRspWrapper = tmqFreeRspWrapper(pRspWrapper);
1965
          pVg->emptyBlockReceiveTs = taosGetTimestampMs();
L
Liu Jicong 已提交
1966
          taosFreeQitem(pollRspWrapper);
1967
        } else {  // build rsp
X
Xiaoyu Wang 已提交
1968
          int64_t    numOfRows = 0;
1969
          SMqRspObj* pRsp = tmqBuildRspFromWrapper(pollRspWrapper, pVg, &numOfRows);
H
Haojun Liao 已提交
1970
          tmq->totalRows += numOfRows;
1971
          pVg->emptyBlockReceiveTs = 0;
H
Haojun Liao 已提交
1972
          tscDebug("consumer:0x%" PRIx64 " process poll rsp, vgId:%d, offset:%s, blocks:%d, rows:%" PRId64
wmmhello's avatar
wmmhello 已提交
1973
                   ", vg total:%" PRId64 ", total:%" PRId64 ", reqId:0x%" PRIx64,
H
Haojun Liao 已提交
1974
                   tmq->consumerId, pVg->vgId, buf, pDataRsp->blockNum, numOfRows, pVg->numOfRows, tmq->totalRows,
H
Haojun Liao 已提交
1975
                   pollRspWrapper->reqId);
1976
          taosFreeQitem(pollRspWrapper);
wmmhello's avatar
wmmhello 已提交
1977
          taosWUnLockLatch(&tmq->lock);
1978 1979
          return pRsp;
        }
wmmhello's avatar
wmmhello 已提交
1980
        taosWUnLockLatch(&tmq->lock);
X
Xiaoyu Wang 已提交
1981
      } 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, pDataRsp->head.epoch, consumerEpoch);
1984
        pRspWrapper = tmqFreeRspWrapper(pRspWrapper);
L
Liu Jicong 已提交
1985 1986
        taosFreeQitem(pollRspWrapper);
      }
1987
    } else if (pRspWrapper->tmqRspType == TMQ_MSG_TYPE__POLL_META_RSP) {
1988
      // todo handle the wal range and epset for each vgroup
1989
      SMqPollRspWrapper* pollRspWrapper = (SMqPollRspWrapper*)pRspWrapper;
L
Liu Jicong 已提交
1990
      int32_t            consumerEpoch = atomic_load_32(&tmq->epoch);
1991 1992 1993

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

L
Liu Jicong 已提交
1994
      if (pollRspWrapper->metaRsp.head.epoch == consumerEpoch) {
wmmhello's avatar
wmmhello 已提交
1995
        taosWLockLatch(&tmq->lock);
wmmhello's avatar
wmmhello 已提交
1996 1997 1998 1999 2000 2001
        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);
wmmhello's avatar
wmmhello 已提交
2002
          taosWUnLockLatch(&tmq->lock);
wmmhello's avatar
wmmhello 已提交
2003 2004
          return NULL;
        }
H
Haojun Liao 已提交
2005

2006
        updateVgInfo(pVg, &pollRspWrapper->metaRsp.rspOffset, &pollRspWrapper->metaRsp.rspOffset, pollRspWrapper->metaRsp.head.walsver, pollRspWrapper->metaRsp.head.walever, tmq->consumerId);
L
Liu Jicong 已提交
2007
        // build rsp
L
Liu Jicong 已提交
2008
        SMqMetaRspObj* pRsp = tmqBuildMetaRspFromWrapper(pollRspWrapper);
L
Liu Jicong 已提交
2009
        taosFreeQitem(pollRspWrapper);
wmmhello's avatar
wmmhello 已提交
2010
        taosWUnLockLatch(&tmq->lock);
L
Liu Jicong 已提交
2011 2012
        return pRsp;
      } else {
H
Haojun Liao 已提交
2013
        tscDebug("consumer:0x%" PRIx64 " vgId:%d msg discard since epoch mismatch: msg epoch %d, consumer epoch %d",
2014
                 tmq->consumerId, pollRspWrapper->vgId, pollRspWrapper->metaRsp.head.epoch, consumerEpoch);
2015
        pRspWrapper = tmqFreeRspWrapper(pRspWrapper);
L
Liu Jicong 已提交
2016
        taosFreeQitem(pollRspWrapper);
X
Xiaoyu Wang 已提交
2017
      }
2018 2019
    } else if (pRspWrapper->tmqRspType == TMQ_MSG_TYPE__TAOSX_RSP) {
      SMqPollRspWrapper* pollRspWrapper = (SMqPollRspWrapper*)pRspWrapper;
X
Xiaoyu Wang 已提交
2020
      int32_t            consumerEpoch = atomic_load_32(&tmq->epoch);
H
Haojun Liao 已提交
2021

L
Liu Jicong 已提交
2022
      if (pollRspWrapper->taosxRsp.head.epoch == consumerEpoch) {
wmmhello's avatar
wmmhello 已提交
2023
        taosWLockLatch(&tmq->lock);
wmmhello's avatar
wmmhello 已提交
2024 2025 2026 2027 2028 2029
        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);
wmmhello's avatar
wmmhello 已提交
2030
          taosWUnLockLatch(&tmq->lock);
wmmhello's avatar
wmmhello 已提交
2031 2032
          return NULL;
        }
H
Haojun Liao 已提交
2033

2034
        updateVgInfo(pVg, &pollRspWrapper->taosxRsp.reqOffset, &pollRspWrapper->taosxRsp.rspOffset, pollRspWrapper->taosxRsp.head.walsver, pollRspWrapper->taosxRsp.head.walever, tmq->consumerId);
H
Haojun Liao 已提交
2035

L
Liu Jicong 已提交
2036
        if (pollRspWrapper->taosxRsp.blockNum == 0) {
wmmhello's avatar
wmmhello 已提交
2037
          tscDebug("consumer:0x%" PRIx64 " taosx empty block received, vgId:%d, vg total:%" PRId64 ", reqId:0x%" PRIx64,
H
Haojun Liao 已提交
2038
                   tmq->consumerId, pVg->vgId, pVg->numOfRows, pollRspWrapper->reqId);
H
Haojun Liao 已提交
2039
          pVg->emptyBlockReceiveTs = taosGetTimestampMs();
2040
          pRspWrapper = tmqFreeRspWrapper(pRspWrapper);
H
Haojun Liao 已提交
2041
          taosFreeQitem(pollRspWrapper);
H
Haojun Liao 已提交
2042
        } else {
X
Xiaoyu Wang 已提交
2043
          pVg->emptyBlockReceiveTs = 0;  // reset the ts
wmmhello's avatar
wmmhello 已提交
2044 2045 2046 2047 2048 2049 2050 2051
          // build rsp
          void*   pRsp = NULL;
          int64_t numOfRows = 0;
          if (pollRspWrapper->taosxRsp.createTableNum == 0) {
            pRsp = tmqBuildRspFromWrapper(pollRspWrapper, pVg, &numOfRows);
          } else {
            pRsp = tmqBuildTaosxRspFromWrapper(pollRspWrapper, pVg, &numOfRows);
          }
2052

wmmhello's avatar
wmmhello 已提交
2053
          tmq->totalRows += numOfRows;
H
Haojun Liao 已提交
2054

wmmhello's avatar
wmmhello 已提交
2055
          char buf[TSDB_OFFSET_LEN] = {0};
wmmhello's avatar
wmmhello 已提交
2056 2057 2058 2059 2060
          tFormatOffset(buf, TSDB_OFFSET_LEN, &pVg->offsetInfo.currentOffset);
          tscDebug("consumer:0x%" PRIx64 " process taosx poll rsp, vgId:%d, offset:%s, blocks:%d, rows:%" PRId64
                       ", vg total:%" PRId64 ", total:%" PRId64 ", reqId:0x%" PRIx64,
                   tmq->consumerId, pVg->vgId, buf, pollRspWrapper->dataRsp.blockNum, numOfRows, pVg->numOfRows,
                   tmq->totalRows, pollRspWrapper->reqId);
H
Haojun Liao 已提交
2061

wmmhello's avatar
wmmhello 已提交
2062 2063 2064 2065 2066
          taosFreeQitem(pollRspWrapper);
          taosWUnLockLatch(&tmq->lock);
          return pRsp;
        }
        taosWUnLockLatch(&tmq->lock);
L
Liu Jicong 已提交
2067
      } else {
H
Haojun Liao 已提交
2068
        tscDebug("consumer:0x%" PRIx64 " vgId:%d msg discard since epoch mismatch: msg epoch %d, consumer epoch %d",
2069
                 tmq->consumerId, pollRspWrapper->vgId, pollRspWrapper->taosxRsp.head.epoch, consumerEpoch);
2070
        pRspWrapper = tmqFreeRspWrapper(pRspWrapper);
L
Liu Jicong 已提交
2071 2072
        taosFreeQitem(pollRspWrapper);
      }
X
Xiaoyu Wang 已提交
2073
    } else {
H
Haojun Liao 已提交
2074 2075
      tscDebug("consumer:0x%" PRIx64 " not data msg received", tmq->consumerId);

X
Xiaoyu Wang 已提交
2076
      bool reset = false;
2077 2078
      tmqHandleNoPollRsp(tmq, pRspWrapper, &reset);
      taosFreeQitem(pRspWrapper);
X
Xiaoyu Wang 已提交
2079
      if (pollIfReset && reset) {
2080
        tscDebug("consumer:0x%" PRIx64 ", reset and repoll", tmq->consumerId);
2081
        tmqPollImpl(tmq, timeout);
X
Xiaoyu Wang 已提交
2082 2083 2084 2085 2086
      }
    }
  }
}

2087
TAOS_RES* tmq_consumer_poll(tmq_t* tmq, int64_t timeout) {
L
Liu Jicong 已提交
2088 2089
  void*   rspObj;
  int64_t startTime = taosGetTimestampMs();
L
Liu Jicong 已提交
2090

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

2094
  // in no topic status, delayed task also need to be processed
L
Liu Jicong 已提交
2095
  if (atomic_load_8(&tmq->status) == TMQ_CONSUMER_STATUS__INIT) {
2096
    tscInfo("consumer:0x%" PRIx64 " poll return since consumer is init", tmq->consumerId);
2097
    taosMsleep(500);  //     sleep for a while
2098 2099 2100
    return NULL;
  }

wmmhello's avatar
wmmhello 已提交
2101
  while (atomic_load_8(&tmq->status) == TMQ_CONSUMER_STATUS__RECOVER) {
L
Liu Jicong 已提交
2102
    int32_t retryCnt = 0;
2103
    while (TSDB_CODE_MND_CONSUMER_NOT_READY == doAskEp(tmq)) {
H
Haojun Liao 已提交
2104
      if (retryCnt++ > 40) {
L
Liu Jicong 已提交
2105 2106
        return NULL;
      }
2107

2108
      tscInfo("consumer:0x%" PRIx64 " not ready, retry:%d/40 in 500ms", tmq->consumerId, retryCnt);
L
Liu Jicong 已提交
2109 2110 2111 2112
      taosMsleep(500);
    }
  }

X
Xiaoyu Wang 已提交
2113
  while (1) {
L
Liu Jicong 已提交
2114
    tmqHandleAllDelayedTask(tmq);
2115

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

2120
    rspObj = tmqHandleAllRsp(tmq, timeout, false);
L
Liu Jicong 已提交
2121
    if (rspObj) {
2122
      tscDebug("consumer:0x%" PRIx64 " return rsp %p", tmq->consumerId, rspObj);
L
Liu Jicong 已提交
2123
      return (TAOS_RES*)rspObj;
L
Liu Jicong 已提交
2124
    } else if (terrno == TSDB_CODE_TQ_NO_COMMITTED_OFFSET) {
2125
      tscInfo("consumer:0x%" PRIx64 " return null since no committed offset", tmq->consumerId);
L
Liu Jicong 已提交
2126
      return NULL;
X
Xiaoyu Wang 已提交
2127
    }
2128

2129
    if (timeout >= 0) {
L
Liu Jicong 已提交
2130
      int64_t currentTime = taosGetTimestampMs();
2131 2132
      int64_t elapsedTime = currentTime - startTime;
      if (elapsedTime > timeout) {
2133
        tscInfo("consumer:0x%" PRIx64 " (epoch %d) timeout, no rsp, start time %" PRId64 ", current time %" PRId64,
L
Liu Jicong 已提交
2134
                 tmq->consumerId, tmq->epoch, startTime, currentTime);
X
Xiaoyu Wang 已提交
2135 2136
        return NULL;
      }
2137
      tsem_timewait(&tmq->rspSem, (timeout - elapsedTime));
L
Liu Jicong 已提交
2138 2139
    } else {
      // use tsem_timewait instead of tsem_wait to avoid unexpected stuck
L
Liu Jicong 已提交
2140
      tsem_timewait(&tmq->rspSem, 1000);
X
Xiaoyu Wang 已提交
2141 2142 2143 2144
    }
  }
}

wmmhello's avatar
wmmhello 已提交
2145 2146
static void displayConsumeStatistics(tmq_t* pTmq) {
  taosRLockLatch(&pTmq->lock);
2147 2148 2149 2150 2151 2152 2153 2154 2155 2156 2157 2158 2159
  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);
2160
    }
2161
  }
wmmhello's avatar
wmmhello 已提交
2162
  taosRUnLockLatch(&pTmq->lock);
2163 2164
  tscDebug("consumer:0x%" PRIx64 " rows dist end", pTmq->consumerId);
}
2165

2166
int32_t tmq_consumer_close(tmq_t* tmq) {
2167
  tscInfo("consumer:0x%" PRIx64 " start to close consumer, status:%d", tmq->consumerId, tmq->status);
2168
  displayConsumeStatistics(tmq);
2169

2170 2171 2172 2173 2174 2175
  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;
2176 2177
      }
    }
2178
    taosSsleep(2);  // sleep 2s for hb to send offset and rows to server
2179

L
Liu Jicong 已提交
2180
    int32_t     retryCnt = 0;
2181
    tmq_list_t* lst = tmq_list_new();
L
Liu Jicong 已提交
2182
    while (1) {
2183
      int32_t rsp = tmq_subscribe(tmq, lst);
L
Liu Jicong 已提交
2184 2185 2186 2187 2188 2189 2190 2191
      if (rsp != TSDB_CODE_MND_CONSUMER_NOT_READY || retryCnt > 5) {
        break;
      } else {
        retryCnt++;
        taosMsleep(500);
      }
    }

2192
    tmq_list_destroy(lst);
2193
  } else {
2194
    tscInfo("consumer:0x%" PRIx64 " not in ready state, close it directly", tmq->consumerId);
L
Liu Jicong 已提交
2195
  }
H
Haojun Liao 已提交
2196

2197
  taosRemoveRef(tmqMgmt.rsetId, tmq->refId);
L
Liu Jicong 已提交
2198
  return 0;
2199
}
L
Liu Jicong 已提交
2200

L
Liu Jicong 已提交
2201 2202
const char* tmq_err2str(int32_t err) {
  if (err == 0) {
L
Liu Jicong 已提交
2203
    return "success";
L
Liu Jicong 已提交
2204
  } else if (err == -1) {
L
Liu Jicong 已提交
2205 2206 2207
    return "fail";
  } else {
    return tstrerror(err);
L
Liu Jicong 已提交
2208 2209
  }
}
L
Liu Jicong 已提交
2210

L
Liu Jicong 已提交
2211 2212 2213 2214 2215
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;
2216 2217
  } else if (TD_RES_TMQ_METADATA(res)) {
    return TMQ_RES_METADATA;
L
Liu Jicong 已提交
2218 2219 2220 2221 2222
  } else {
    return TMQ_RES_INVALID;
  }
}

L
Liu Jicong 已提交
2223
const char* tmq_get_topic_name(TAOS_RES* res) {
L
Liu Jicong 已提交
2224 2225
  if (TD_RES_TMQ(res)) {
    SMqRspObj* pRspObj = (SMqRspObj*)res;
L
Liu Jicong 已提交
2226
    return strchr(pRspObj->topic, '.') + 1;
L
Liu Jicong 已提交
2227 2228 2229
  } else if (TD_RES_TMQ_META(res)) {
    SMqMetaRspObj* pMetaRspObj = (SMqMetaRspObj*)res;
    return strchr(pMetaRspObj->topic, '.') + 1;
2230 2231 2232
  } else if (TD_RES_TMQ_METADATA(res)) {
    SMqTaosxRspObj* pRspObj = (SMqTaosxRspObj*)res;
    return strchr(pRspObj->topic, '.') + 1;
L
Liu Jicong 已提交
2233 2234 2235 2236 2237
  } else {
    return NULL;
  }
}

L
Liu Jicong 已提交
2238 2239 2240 2241
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 已提交
2242 2243 2244
  } else if (TD_RES_TMQ_META(res)) {
    SMqMetaRspObj* pMetaRspObj = (SMqMetaRspObj*)res;
    return strchr(pMetaRspObj->db, '.') + 1;
2245 2246 2247
  } else if (TD_RES_TMQ_METADATA(res)) {
    SMqTaosxRspObj* pRspObj = (SMqTaosxRspObj*)res;
    return strchr(pRspObj->db, '.') + 1;
L
Liu Jicong 已提交
2248 2249 2250 2251 2252
  } else {
    return NULL;
  }
}

L
Liu Jicong 已提交
2253 2254 2255 2256
int32_t tmq_get_vgroup_id(TAOS_RES* res) {
  if (TD_RES_TMQ(res)) {
    SMqRspObj* pRspObj = (SMqRspObj*)res;
    return pRspObj->vgId;
L
Liu Jicong 已提交
2257 2258 2259
  } else if (TD_RES_TMQ_META(res)) {
    SMqMetaRspObj* pMetaRspObj = (SMqMetaRspObj*)res;
    return pMetaRspObj->vgId;
2260
  } else if (TD_RES_TMQ_METADATA(res)) {
2261 2262
    SMqTaosxRspObj* pRspObj = (SMqTaosxRspObj*)res;
    return pRspObj->vgId;
L
Liu Jicong 已提交
2263 2264 2265 2266
  } else {
    return -1;
  }
}
L
Liu Jicong 已提交
2267

2268 2269 2270
int64_t tmq_get_vgroup_offset(TAOS_RES* res) {
  if (TD_RES_TMQ(res)) {
    SMqRspObj* pRspObj = (SMqRspObj*) res;
2271
    STqOffsetVal* pOffset = &pRspObj->rsp.reqOffset;
2272
    if (pOffset->type == TMQ_OFFSET__LOG) {
2273
      return pRspObj->rsp.reqOffset.version;
2274 2275
    }else{
      tscError("invalid offset type:%d", pOffset->type);
2276 2277 2278 2279 2280 2281 2282 2283
    }
  } 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;
2284 2285
    if (pRspObj->rsp.reqOffset.type == TMQ_OFFSET__LOG) {
      return pRspObj->rsp.reqOffset.version;
2286
    }
2287
  } else{
2288
    tscError("invalid tmq type:%d", *(int8_t*)res);
2289 2290 2291 2292 2293 2294
  }

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

L
Liu Jicong 已提交
2295 2296 2297 2298 2299 2300 2301
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;
    }
2302
    return (const char*)taosArrayGetP(pRspObj->rsp.blockTbName, pRspObj->resIter);
2303 2304
  } else if (TD_RES_TMQ_METADATA(res)) {
    SMqTaosxRspObj* pRspObj = (SMqTaosxRspObj*)res;
L
Liu Jicong 已提交
2305 2306 2307
    if (!pRspObj->rsp.withTbName || pRspObj->rsp.blockTbName == NULL || pRspObj->resIter < 0 ||
        pRspObj->resIter >= pRspObj->rsp.blockNum) {
      return NULL;
2308
    }
L
Liu Jicong 已提交
2309 2310
    return (const char*)taosArrayGetP(pRspObj->rsp.blockTbName, pRspObj->resIter);
  }
L
Liu Jicong 已提交
2311 2312
  return NULL;
}
2313

2314 2315 2316 2317
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
2318
    asyncCommitOffset(tmq, pRes, TDMT_VND_TMQ_COMMIT_OFFSET, cb, param);
2319
  }
L
Liu Jicong 已提交
2320 2321
}

2322
static void commitCallBackFn(tmq_t *UNUSED_PARAM(tmq), int32_t code, void* param) {
2323 2324 2325
  SSyncCommitInfo* pInfo = (SSyncCommitInfo*) param;
  pInfo->code = code;
  tsem_post(&pInfo->sem);
2326
}
2327

2328 2329 2330 2331 2332 2333 2334 2335 2336
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 已提交
2337
  } else {
2338
    asyncCommitOffset(tmq, pRes, TDMT_VND_TMQ_COMMIT_OFFSET, commitCallBackFn, pInfo);
2339 2340
  }

2341 2342
  tsem_wait(&pInfo->sem);
  code = pInfo->code;
H
Haojun Liao 已提交
2343 2344

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

X
Xiaoyu Wang 已提交
2347
  tscDebug("consumer:0x%" PRIx64 " sync commit done, code:%s", tmq->consumerId, tstrerror(code));
2348 2349 2350 2351 2352 2353 2354 2355 2356 2357 2358 2359
  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);
2360
    doUpdateLocalEp(pTmq, head->epoch, &rsp);
2361 2362 2363
    tDeleteSMqAskEpRsp(&rsp);
  }

H
Haojun Liao 已提交
2364
  tsem_post(&pInfo->sem);
2365 2366 2367 2368 2369 2370 2371 2372 2373 2374 2375 2376 2377 2378 2379 2380 2381 2382 2383 2384 2385 2386 2387 2388 2389 2390
}

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 已提交
2391
  tsem_init(&pInfo->sem, 0, 0);
2392 2393

  asyncAskEp(pTmq, updateEpCallbackFn, pInfo);
H
Haojun Liao 已提交
2394
  tsem_wait(&pInfo->sem);
2395 2396

  int32_t code = pInfo->code;
H
Haojun Liao 已提交
2397
  tsem_destroy(&pInfo->sem);
2398 2399 2400 2401 2402
  taosMemoryFree(pInfo);
  return code;
}

void asyncAskEp(tmq_t* pTmq, __tmq_askep_fn_t askEpFn, void* param) {
2403
  SMqAskEpReq req = {0};
2404 2405 2406
  req.consumerId = pTmq->consumerId;
  req.epoch = pTmq->epoch;
  strcpy(req.cgroup, pTmq->groupId);
2407 2408 2409

  int32_t tlen = tSerializeSMqAskEpReq(NULL, 0, &req);
  if (tlen < 0) {
2410 2411 2412
    tscError("consumer:0x%" PRIx64 ", tSerializeSMqAskEpReq failed", pTmq->consumerId);
    askEpFn(pTmq, TSDB_CODE_INVALID_PARA, NULL, param);
    return;
2413 2414 2415 2416
  }

  void* pReq = taosMemoryCalloc(1, tlen);
  if (pReq == NULL) {
2417 2418 2419
    tscError("consumer:0x%" PRIx64 ", failed to malloc askEpReq msg, size:%d", pTmq->consumerId, tlen);
    askEpFn(pTmq, TSDB_CODE_OUT_OF_MEMORY, NULL, param);
    return;
2420 2421 2422
  }

  if (tSerializeSMqAskEpReq(pReq, tlen, &req) < 0) {
2423
    tscError("consumer:0x%" PRIx64 ", tSerializeSMqAskEpReq %d failed", pTmq->consumerId, tlen);
2424
    taosMemoryFree(pReq);
2425 2426 2427

    askEpFn(pTmq, TSDB_CODE_INVALID_PARA, NULL, param);
    return;
2428 2429 2430 2431
  }

  SMqAskEpCbParam* pParam = taosMemoryCalloc(1, sizeof(SMqAskEpCbParam));
  if (pParam == NULL) {
2432
    tscError("consumer:0x%" PRIx64 ", failed to malloc subscribe param", pTmq->consumerId);
2433
    taosMemoryFree(pReq);
2434 2435 2436

    askEpFn(pTmq, TSDB_CODE_OUT_OF_MEMORY, NULL, param);
    return;
2437 2438
  }

2439 2440 2441 2442
  pParam->refId = pTmq->refId;
  pParam->epoch = pTmq->epoch;
  pParam->pUserFn = askEpFn;
  pParam->pParam = param;
2443 2444 2445 2446 2447

  SMsgSendInfo* sendInfo = taosMemoryCalloc(1, sizeof(SMsgSendInfo));
  if (sendInfo == NULL) {
    taosMemoryFree(pParam);
    taosMemoryFree(pReq);
2448 2449
    askEpFn(pTmq, TSDB_CODE_OUT_OF_MEMORY, NULL, param);
    return;
2450 2451
  }

X
Xiaoyu Wang 已提交
2452
  sendInfo->msgInfo = (SDataBuf){.pData = pReq, .len = tlen, .handle = NULL};
2453 2454 2455 2456

  sendInfo->requestId = generateRequestId();
  sendInfo->requestObjRefId = 0;
  sendInfo->param = pParam;
2457
  sendInfo->fp = askEpCallbackFn;
2458 2459
  sendInfo->msgType = TDMT_MND_TMQ_ASK_EP;

2460
  SEpSet epSet = getEpSet_s(&pTmq->pTscObj->pAppInfo->mgmtEp);
2461
  tscInfo("consumer:0x%" PRIx64 " ask ep from mnode, reqId:0x%" PRIx64, pTmq->consumerId, sendInfo->requestId);
2462 2463

  int64_t transporterId = 0;
2464
  asyncSendMsgToServer(pTmq->pTscObj->pAppInfo->pTransporter, &epSet, &transporterId, sendInfo);
2465 2466 2467 2468 2469 2470 2471
}

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

int32_t tmqCommitDone(SMqCommitCbParamSet* pParamSet) {
2472 2473 2474
  int64_t refId = pParamSet->refId;

  tmq_t* tmq = taosAcquireRef(tmqMgmt.rsetId, refId);
2475 2476 2477 2478 2479 2480 2481
  if (tmq == NULL) {
    taosMemoryFree(pParamSet);
    terrno = TSDB_CODE_TMQ_CONSUMER_CLOSED;
    return -1;
  }

  // if no more waiting rsp
2482
  pParamSet->callbackFn(tmq, pParamSet->code, pParamSet->userParam);
2483
  taosMemoryFree(pParamSet);
wmmhello's avatar
wmmhello 已提交
2484
//  tmq->needReportOffsetRows = true;
2485 2486

  taosReleaseRef(tmqMgmt.rsetId, refId);
2487
  return 0;
2488 2489
}

2490
void commitRspCountDown(SMqCommitCbParamSet* pParamSet, int64_t consumerId, const char* pTopic, int32_t vgId) {
2491 2492
  int32_t waitingRspNum = atomic_sub_fetch_32(&pParamSet->waitingRspNum, 1);
  if (waitingRspNum == 0) {
H
Haojun Liao 已提交
2493 2494
    tscDebug("consumer:0x%" PRIx64 " topic:%s vgId:%d all commit-rsp received, commit completed", consumerId, pTopic,
             vgId);
2495
    tmqCommitDone(pParamSet);
H
Haojun Liao 已提交
2496 2497 2498
  } else {
    tscDebug("consumer:0x%" PRIx64 " topic:%s vgId:%d commit-rsp received, remain:%d", consumerId, pTopic, vgId,
             waitingRspNum);
2499 2500
  }
}
2501 2502 2503 2504 2505 2506 2507 2508 2509 2510 2511 2512 2513 2514 2515 2516 2517 2518 2519 2520 2521 2522

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 已提交
2523 2524
}

2525 2526 2527 2528 2529 2530 2531 2532 2533 2534 2535 2536 2537 2538 2539 2540 2541 2542 2543 2544 2545
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,
2546
                                       .vgId = pParam->vgId};
2547 2548 2549 2550 2551 2552 2553 2554 2555 2556

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

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

2557 2558
  taosMemoryFree(pMsg->pData);
  taosMemoryFree(pMsg->pEpSet);
2559 2560 2561 2562 2563
  taosMemoryFree(pParam);
  return 0;
}

static void destroyCommonInfo(SMqVgCommon* pCommon) {
wmmhello's avatar
wmmhello 已提交
2564 2565 2566
  if(pCommon == NULL){
    return;
  }
2567 2568 2569 2570 2571 2572 2573
  taosArrayDestroy(pCommon->pList);
  tsem_destroy(&pCommon->rsp);
  taosThreadMutexDestroy(&pCommon->mutex);
  taosMemoryFree(pCommon->pTopicName);
  taosMemoryFree(pCommon);
}

2574 2575 2576 2577 2578 2579 2580
static bool isInSnapshotMode(int8_t type, bool useSnapshot){
  if ((type < TMQ_OFFSET__LOG && useSnapshot) || type > TMQ_OFFSET__LOG) {
    return true;
  }
  return false;
}

H
Haojun Liao 已提交
2581
int32_t tmq_get_topic_assignment(tmq_t* tmq, const char* pTopicName, tmq_topic_assignment** assignment,
H
Haojun Liao 已提交
2582
                                 int32_t* numOfAssignment) {
H
Haojun Liao 已提交
2583 2584
  *numOfAssignment = 0;
  *assignment = NULL;
wmmhello's avatar
wmmhello 已提交
2585
  SMqVgCommon* pCommon = NULL;
H
Haojun Liao 已提交
2586

2587
  int32_t accId = tmq->pTscObj->acctId;
2588
  char    tname[128] = {0};
2589
  sprintf(tname, "%d.%s", accId, pTopicName);
wmmhello's avatar
wmmhello 已提交
2590
  int32_t code = TSDB_CODE_SUCCESS;
2591

wmmhello's avatar
wmmhello 已提交
2592
  taosWLockLatch(&tmq->lock);
2593
  SMqClientTopic* pTopic = getTopicByName(tmq, tname);
H
Haojun Liao 已提交
2594
  if (pTopic == NULL) {
wmmhello's avatar
wmmhello 已提交
2595 2596
    code = TSDB_CODE_INVALID_PARA;
    goto end;
H
Haojun Liao 已提交
2597 2598 2599 2600
  }

  // in case of snapshot is opened, no valid offset will return
  *numOfAssignment = taosArrayGetSize(pTopic->vgs);
2601 2602
  for (int32_t j = 0; j < (*numOfAssignment); ++j) {
    SMqClientVg* pClientVg = taosArrayGet(pTopic->vgs, j);
2603 2604 2605
    int32_t type = pClientVg->offsetInfo.currentOffset.type;
    if (isInSnapshotMode(type, tmq->useSnapshot)) {
      tscError("consumer:0x%" PRIx64 " offset type:%d not wal version, assignment not allowed", tmq->consumerId, type);
2606 2607 2608 2609
      code = TSDB_CODE_TMQ_SNAPSHOT_ERROR;
      goto end;
    }
  }
2610 2611 2612 2613 2614

  *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));
wmmhello's avatar
wmmhello 已提交
2615 2616
    code = TSDB_CODE_OUT_OF_MEMORY;
    goto end;
2617 2618
  }

2619 2620
  bool needFetch = false;

H
Haojun Liao 已提交
2621 2622
  for (int32_t j = 0; j < (*numOfAssignment); ++j) {
    SMqClientVg* pClientVg = taosArrayGet(pTopic->vgs, j);
2623
    if (pClientVg->offsetInfo.currentOffset.type != TMQ_OFFSET__LOG) {
2624 2625 2626
      needFetch = true;
      break;
    }
H
Haojun Liao 已提交
2627 2628

    tmq_topic_assignment* pAssignment = &(*assignment)[j];
2629
    pAssignment->currentOffset = pClientVg->offsetInfo.seekOffset.version;
H
Haojun Liao 已提交
2630 2631
    pAssignment->begin = pClientVg->offsetInfo.walVerBegin;
    pAssignment->end = pClientVg->offsetInfo.walVerEnd;
2632
    pAssignment->vgId = pClientVg->vgId;
wmmhello's avatar
wmmhello 已提交
2633 2634
    tscInfo("consumer:0x%" PRIx64 " get assignment from local:%d->%" PRId64, tmq->consumerId,
            pAssignment->vgId, pAssignment->currentOffset);
H
Haojun Liao 已提交
2635 2636
  }

2637
  if (needFetch) {
wmmhello's avatar
wmmhello 已提交
2638
    pCommon = taosMemoryCalloc(1, sizeof(SMqVgCommon));
2639 2640
    if (pCommon == NULL) {
      terrno = TSDB_CODE_OUT_OF_MEMORY;
wmmhello's avatar
wmmhello 已提交
2641 2642
      code = terrno;
      goto end;
2643 2644 2645 2646 2647 2648 2649 2650 2651 2652 2653 2654 2655 2656
    }

    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) {
wmmhello's avatar
wmmhello 已提交
2657 2658
        code = terrno;
        goto end;
2659 2660 2661 2662 2663 2664 2665 2666 2667
      }

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

      SMqPollReq req = {0};
      tmqBuildConsumeReqImpl(&req, tmq, 10, pTopic, pClientVg);
2668
      req.reqOffset = pClientVg->offsetInfo.seekOffset;
2669 2670 2671 2672

      int32_t msgSize = tSerializeSMqPollReq(NULL, 0, &req);
      if (msgSize < 0) {
        taosMemoryFree(pParam);
wmmhello's avatar
wmmhello 已提交
2673 2674
        code = terrno;
        goto end;
2675 2676 2677 2678 2679
      }

      char* msg = taosMemoryCalloc(1, msgSize);
      if (NULL == msg) {
        taosMemoryFree(pParam);
wmmhello's avatar
wmmhello 已提交
2680 2681
        code = terrno;
        goto end;
2682 2683 2684 2685 2686
      }

      if (tSerializeSMqPollReq(msg, msgSize, &req) < 0) {
        taosMemoryFree(msg);
        taosMemoryFree(pParam);
wmmhello's avatar
wmmhello 已提交
2687 2688
        code = terrno;
        goto end;
2689 2690 2691 2692 2693 2694
      }

      SMsgSendInfo* sendInfo = taosMemoryCalloc(1, sizeof(SMsgSendInfo));
      if (sendInfo == NULL) {
        taosMemoryFree(pParam);
        taosMemoryFree(msg);
wmmhello's avatar
wmmhello 已提交
2695 2696
        code = terrno;
        goto end;
2697 2698 2699 2700 2701 2702 2703 2704 2705 2706
      }

      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;
wmmhello's avatar
wmmhello 已提交
2707
      char    offsetFormatBuf[TSDB_OFFSET_LEN] = {0};
2708 2709
      tFormatOffset(offsetFormatBuf, tListLen(offsetFormatBuf), &pClientVg->offsetInfo.currentOffset);

2710
      tscInfo("consumer:0x%" PRIx64 " %s retrieve wal info vgId:%d, epoch %d, req:%s, reqId:0x%" PRIx64,
wmmhello's avatar
wmmhello 已提交
2711
              tmq->consumerId, pTopic->topicName, pClientVg->vgId, tmq->epoch, offsetFormatBuf, req.reqId);
2712 2713 2714 2715
      asyncSendMsgToServer(tmq->pTscObj->pAppInfo->pTransporter, &pClientVg->epSet, &transporterId, sendInfo);
    }

    tsem_wait(&pCommon->rsp);
wmmhello's avatar
wmmhello 已提交
2716
    code = pCommon->code;
2717 2718 2719

    terrno = code;
    if (code != TSDB_CODE_SUCCESS) {
wmmhello's avatar
wmmhello 已提交
2720
      goto end;
2721
    }
wmmhello's avatar
wmmhello 已提交
2722 2723 2724
    int32_t num = taosArrayGetSize(pCommon->pList);
    for(int32_t i = 0; i < num; ++i) {
      (*assignment)[i] = *(tmq_topic_assignment*)taosArrayGet(pCommon->pList, i);
2725
    }
wmmhello's avatar
wmmhello 已提交
2726
    *numOfAssignment = num;
2727

2728 2729 2730 2731 2732 2733 2734 2735 2736 2737
    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;
wmmhello's avatar
wmmhello 已提交
2738
        tscInfo("vgId:%d offset is update to:%"PRId64, p->vgId, p->currentOffset);
2739 2740 2741 2742 2743

        pOffsetInfo->walVerBegin = p->begin;
        pOffsetInfo->walVerEnd = p->end;
      }
    }
wmmhello's avatar
wmmhello 已提交
2744
  }
2745

wmmhello's avatar
wmmhello 已提交
2746 2747 2748 2749 2750
end:
  if(code != TSDB_CODE_SUCCESS){
    taosMemoryFree(*assignment);
    *assignment = NULL;
    *numOfAssignment = 0;
2751
  }
wmmhello's avatar
wmmhello 已提交
2752 2753 2754
  destroyCommonInfo(pCommon);
  taosWUnLockLatch(&tmq->lock);
  return code;
H
Haojun Liao 已提交
2755 2756
}

T
t_max 已提交
2757 2758 2759 2760 2761 2762 2763 2764
void tmq_free_assignment(tmq_topic_assignment* pAssignment) {
    if (pAssignment == NULL) {
        return;
    }

    taosMemoryFree(pAssignment);
}

2765 2766 2767 2768 2769 2770 2771 2772 2773 2774 2775
static int32_t tmqSeekCb(void* param, SDataBuf* pMsg, int32_t code) {
  if (pMsg) {
    taosMemoryFree(pMsg->pData);
    taosMemoryFree(pMsg->pEpSet);
  }
  SMqSeekParam* pParam = param;
  pParam->code = code;
  tsem_post(&pParam->sem);
  return 0;
}

2776
int32_t tmq_offset_seek(tmq_t* tmq, const char* pTopicName, int32_t vgId, int64_t offset) {
H
Haojun Liao 已提交
2777
  if (tmq == NULL) {
H
Haojun Liao 已提交
2778
    tscError("invalid tmq handle, null");
H
Haojun Liao 已提交
2779 2780 2781
    return TSDB_CODE_INVALID_PARA;
  }

2782 2783 2784 2785
  int32_t accId = tmq->pTscObj->acctId;
  char tname[128] = {0};
  sprintf(tname, "%d.%s", accId, pTopicName);

wmmhello's avatar
wmmhello 已提交
2786
  taosWLockLatch(&tmq->lock);
2787
  SMqClientTopic* pTopic = getTopicByName(tmq, tname);
H
Haojun Liao 已提交
2788
  if (pTopic == NULL) {
2789
    tscError("consumer:0x%" PRIx64 " invalid topic name:%s", tmq->consumerId, pTopicName);
wmmhello's avatar
wmmhello 已提交
2790
    taosWUnLockLatch(&tmq->lock);
2791
    return TSDB_CODE_TMQ_INVALID_TOPIC;
H
Haojun Liao 已提交
2792 2793 2794
  }

  SMqClientVg* pVg = NULL;
H
Haojun Liao 已提交
2795 2796
  int32_t      numOfVgs = taosArrayGetSize(pTopic->vgs);
  for (int32_t i = 0; i < numOfVgs; ++i) {
H
Haojun Liao 已提交
2797
    SMqClientVg* pClientVg = taosArrayGet(pTopic->vgs, i);
2798
    if (pClientVg->vgId == vgId) {
H
Haojun Liao 已提交
2799 2800 2801 2802 2803 2804
      pVg = pClientVg;
      break;
    }
  }

  if (pVg == NULL) {
2805
    tscError("consumer:0x%" PRIx64 " invalid vgroup id:%d", tmq->consumerId, vgId);
wmmhello's avatar
wmmhello 已提交
2806
    taosWUnLockLatch(&tmq->lock);
2807
    return TSDB_CODE_TMQ_INVALID_VGID;
H
Haojun Liao 已提交
2808 2809 2810 2811 2812
  }

  SVgOffsetInfo* pOffsetInfo = &pVg->offsetInfo;

  int32_t type = pOffsetInfo->currentOffset.type;
2813
  if (isInSnapshotMode(type, tmq->useSnapshot)) {
wmmhello's avatar
wmmhello 已提交
2814 2815
    tscError("consumer:0x%" PRIx64 " offset type:%d not wal version, seek not allowed", tmq->consumerId, type);
    taosWUnLockLatch(&tmq->lock);
2816
    return TSDB_CODE_TMQ_SNAPSHOT_ERROR;
wmmhello's avatar
wmmhello 已提交
2817 2818 2819 2820 2821 2822
  }

  if (type == TMQ_OFFSET__LOG && (offset < pOffsetInfo->walVerBegin || offset > pOffsetInfo->walVerEnd)) {
    tscError("consumer:0x%" PRIx64 " invalid seek params, offset:%" PRId64 ", valid range:[%" PRId64 ", %" PRId64 "]",
             tmq->consumerId, offset, pOffsetInfo->walVerBegin, pOffsetInfo->walVerEnd);
    taosWUnLockLatch(&tmq->lock);
2823
    return TSDB_CODE_TMQ_VERSION_OUT_OF_RANGE;
wmmhello's avatar
wmmhello 已提交
2824
  }
H
Haojun Liao 已提交
2825

H
Haojun Liao 已提交
2826
  // update the offset, and then commit to vnode
wmmhello's avatar
wmmhello 已提交
2827 2828
  pOffsetInfo->currentOffset.type = TMQ_OFFSET__LOG;
  pOffsetInfo->currentOffset.version = offset >= 1 ? offset - 1 : 0;
2829
  pOffsetInfo->seekOffset = pOffsetInfo->currentOffset;
2830
//  pOffsetInfo->committedOffset.version = INT64_MIN;
wmmhello's avatar
wmmhello 已提交
2831
  pVg->seekUpdated = true;
2832 2833 2834 2835 2836 2837 2838 2839 2840 2841 2842 2843 2844 2845 2846 2847 2848 2849 2850 2851 2852 2853 2854 2855 2856 2857 2858 2859 2860 2861 2862 2863 2864 2865 2866 2867 2868 2869 2870 2871 2872 2873 2874 2875 2876 2877 2878
  tscInfo("consumer:0x%" PRIx64 " seek to %" PRId64 " on vgId:%d", tmq->consumerId, offset, vgId);

  SMqSeekReq req = {0};
  snprintf(req.subKey, TSDB_SUBSCRIBE_KEY_LEN, "%s:%s", tmq->groupId, pTopic->topicName);
  req.head.vgId = pVg->vgId;
  req.consumerId = tmq->consumerId;

  int32_t msgSize = tSerializeSMqSeekReq(NULL, 0, &req);
  if (msgSize < 0) {
    taosWUnLockLatch(&tmq->lock);
    return TSDB_CODE_PAR_INTERNAL_ERROR;
  }

  char* msg = taosMemoryCalloc(1, msgSize);
  if (NULL == msg) {
    taosWUnLockLatch(&tmq->lock);
    return TSDB_CODE_OUT_OF_MEMORY;
  }

  if (tSerializeSMqSeekReq(msg, msgSize, &req) < 0) {
    taosMemoryFree(msg);
    taosWUnLockLatch(&tmq->lock);
    return TSDB_CODE_PAR_INTERNAL_ERROR;
  }

  SMsgSendInfo* sendInfo = taosMemoryCalloc(1, sizeof(SMsgSendInfo));
  if (sendInfo == NULL) {
    taosMemoryFree(msg);
    taosWUnLockLatch(&tmq->lock);
    return TSDB_CODE_OUT_OF_MEMORY;
  }

  SMqSeekParam* pParam = taosMemoryMalloc(sizeof(SMqSeekParam));
  if (pParam == NULL) {
    taosMemoryFree(msg);
    taosMemoryFree(sendInfo);
    taosWUnLockLatch(&tmq->lock);
    return TSDB_CODE_OUT_OF_MEMORY;
  }
  tsem_init(&pParam->sem, 0, 0);

  sendInfo->msgInfo = (SDataBuf){.pData = msg, .len = msgSize, .handle = NULL};
  sendInfo->requestId = generateRequestId();
  sendInfo->requestObjRefId = 0;
  sendInfo->param = pParam;
  sendInfo->fp = tmqSeekCb;
  sendInfo->msgType = TDMT_VND_TMQ_SEEK;
H
Haojun Liao 已提交
2879

2880 2881 2882 2883
  int64_t transporterId = 0;
  tscInfo("consumer:0x%" PRIx64 " %s send seek info vgId:%d, epoch %d" PRIx64,
          tmq->consumerId, pTopic->topicName, vgId, tmq->epoch);
  asyncSendMsgToServer(tmq->pTscObj->pAppInfo->pTransporter, &pVg->epSet, &transporterId, sendInfo);
wmmhello's avatar
wmmhello 已提交
2884
  taosWUnLockLatch(&tmq->lock);
2885

2886 2887 2888 2889
  tsem_wait(&pParam->sem);
  int32_t code = pParam->code;
  tsem_destroy(&pParam->sem);
  taosMemoryFree(pParam);
H
Haojun Liao 已提交
2890

2891 2892 2893 2894 2895
  if (code != TSDB_CODE_SUCCESS) {
    tscError("consumer:0x%" PRIx64 " failed to send seek to vgId:%d, code:%s", tmq->consumerId, vgId, tstrerror(code));
  }

  return code;
2896
}