• Home
  • Features
  • Pricing
  • Docs
  • Announcements
  • Sign In

taosdata / TDengine / #4892

20 Dec 2025 01:15PM UTC coverage: 65.571% (+0.02%) from 65.549%
#4892

push

travis-ci

web-flow
feat: support taos_connect_with func (#33952)

* feat: support taos_connect_with

* refactor: enhance connection options and add tests for taos_set_option and taos_connect_with

* fix: handle NULL keys and values in taos_connect_with options

* fix: revert TAOSWS_GIT_TAG to default value "main"

* docs: add TLS configuration options for WebSocket connections in documentation

* docs: modify zh docs and add en docs

* chore: update taos.cfg

* docs: add examples

* docs: add error handling for connection failure in example code

2 of 82 new or added lines in 3 files covered. (2.44%)

527 existing lines in 120 files now uncovered.

182859 of 278870 relevant lines covered (65.57%)

104634355.9 hits per line

Source File
Press 'n' to go to next uncovered line, 'b' for previous

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

16
#include "function.h"
17
#include "functionMgt.h"
18
#include "index.h"
19
#include "os.h"
20
#include "query.h"
21
#include "querynodes.h"
22
#include "taoserror.h"
23
#include "tarray.h"
24
#include "tcompare.h"
25
#include "tdatablock.h"
26
#include "thash.h"
27
#include "tmsg.h"
28
#include "ttime.h"
29

30
#include "executil.h"
31
#include "executorInt.h"
32
#include "querytask.h"
33
#include "storageapi.h"
34
#include "tutil.h"
35
#include "tjson.h"
36

37
typedef struct tagFilterAssist {
38
  SHashObj* colHash;
39
  int32_t   index;
40
  SArray*   cInfoList;
41
  int32_t   code;
42
} tagFilterAssist;
43

44
typedef struct STransTagExprCtx {
45
  int32_t      code;
46
  SMetaReader* pReader;
47
} STransTagExprCtx;
48

49
typedef enum {
50
  FILTER_NO_LOGIC = 1,
51
  FILTER_AND,
52
  FILTER_OTHER,
53
} FilterCondType;
54

55
static FilterCondType checkTagCond(SNode* cond);
56
static int32_t optimizeTbnameInCond(void* metaHandle, int64_t suid, SArray* list, SNode* pTagCond, SStorageAPI* pAPI);
57
static int32_t optimizeTbnameInCondImpl(void* metaHandle, SArray* list, SNode* pTagCond, SStorageAPI* pStoreAPI,
58
                                        uint64_t suid);
59

60
static int32_t getTableList(void* pVnode, SScanPhysiNode* pScanNode, SNode* pTagCond, SNode* pTagIndexCond,
61
                            STableListInfo* pListInfo, uint8_t* digest, const char* idstr, SStorageAPI* pStorageAPI, void* pStreamInfo);
62

63
static int64_t getLimit(const SNode* pLimit) {
562,070,954✔
64
  return (NULL == pLimit || NULL == ((SLimitNode*)pLimit)->limit) ? -1 : ((SLimitNode*)pLimit)->limit->datum.i;
562,070,954✔
65
}
66
static int64_t getOffset(const SNode* pLimit) {
562,075,934✔
67
  return (NULL == pLimit || NULL == ((SLimitNode*)pLimit)->offset) ? -1 : ((SLimitNode*)pLimit)->offset->datum.i;
562,075,934✔
68
}
69
static void releaseColInfoData(void* pCol);
70

71
void initResultRowInfo(SResultRowInfo* pResultRowInfo) {
220,105,567✔
72
  pResultRowInfo->size = 0;
220,105,567✔
73
  pResultRowInfo->cur.pageId = -1;
220,144,925✔
74
}
220,174,823✔
75

76
void closeResultRow(SResultRow* pResultRow) { pResultRow->closed = true; }
5,986,196✔
77

78
void resetResultRow(SResultRow* pResultRow, size_t entrySize) {
405,919,783✔
79
  pResultRow->numOfRows = 0;
405,919,783✔
80
  pResultRow->closed = false;
405,922,063✔
81
  pResultRow->endInterp = false;
405,919,023✔
82
  pResultRow->startInterp = false;
405,920,353✔
83

84
  if (entrySize > 0) {
405,926,433✔
85
    memset(pResultRow->pEntryInfo, 0, entrySize);
405,926,623✔
86
  }
87
}
405,927,003✔
88

89
// TODO refactor: use macro
90
SResultRowEntryInfo* getResultEntryInfo(const SResultRow* pRow, int32_t index, const int32_t* offset) {
2,147,483,647✔
91
  return (SResultRowEntryInfo*)((char*)pRow->pEntryInfo + offset[index]);
2,147,483,647✔
92
}
93

94
size_t getResultRowSize(SqlFunctionCtx* pCtx, int32_t numOfOutput) {
129,695,653✔
95
  int32_t rowSize = (numOfOutput * sizeof(SResultRowEntryInfo)) + sizeof(SResultRow);
129,695,653✔
96

97
  for (int32_t i = 0; i < numOfOutput; ++i) {
644,413,238✔
98
    rowSize += pCtx[i].resDataInfo.interBufSize;
514,725,270✔
99
  }
100

101
  return rowSize;
129,687,968✔
102
}
103

104
// Convert buf read from rocksdb to result row
105
int32_t getResultRowFromBuf(SExprSupp* pSup, const char* inBuf, size_t inBufSize, char** outBuf, size_t* outBufSize) {
×
106
  if (inBuf == NULL || pSup == NULL) {
×
107
    qError("invalid input parameters, inBuf:%p, pSup:%p", inBuf, pSup);
×
108
    return TSDB_CODE_INVALID_PARA;
×
109
  }
110
  SqlFunctionCtx* pCtx = pSup->pCtx;
×
111
  int32_t*        offset = pSup->rowEntryInfoOffset;
×
112
  SResultRow*     pResultRow = NULL;
×
113
  size_t          processedSize = 0;
×
114
  int32_t         code = TSDB_CODE_SUCCESS;
×
115

116
  // calculate the size of output buffer
117
  *outBufSize = getResultRowSize(pCtx, pSup->numOfExprs);
×
118
  *outBuf = taosMemoryMalloc(*outBufSize);
×
119
  if (*outBuf == NULL) {
×
120
    qError("failed to allocate memory for output buffer, size:%zu", *outBufSize);
×
121
    return terrno;
×
122
  }
123
  pResultRow = (SResultRow*)*outBuf;
×
124
  (void)memcpy(pResultRow, inBuf, sizeof(SResultRow));
×
125
  inBuf += sizeof(SResultRow);
×
126
  processedSize += sizeof(SResultRow);
×
127

128
  for (int32_t i = 0; i < pSup->numOfExprs; ++i) {
×
129
    int32_t len = *(int32_t*)inBuf;
×
130
    inBuf += sizeof(int32_t);
×
131
    processedSize += sizeof(int32_t);
×
132
    if (pResultRow->version != FUNCTION_RESULT_INFO_VERSION && pCtx->fpSet.decode) {
×
133
      code = pCtx->fpSet.decode(&pCtx[i], inBuf, getResultEntryInfo(pResultRow, i, offset), pResultRow->version);
×
134
      if (code != TSDB_CODE_SUCCESS) {
×
135
        qError("failed to decode result row, code:%d", code);
×
136
        return code;
×
137
      }
138
    } else {
139
      (void)memcpy(getResultEntryInfo(pResultRow, i, offset), inBuf, len);
×
140
    }
141
    inBuf += len;
×
142
    processedSize += len;
×
143
  }
144

145
  if (processedSize < inBufSize) {
×
146
    // stream stores extra data after result row
147
    size_t leftLen = inBufSize - processedSize;
×
148
    TAOS_MEMORY_REALLOC(*outBuf, *outBufSize + leftLen);
×
149
    if (*outBuf == NULL) {
×
150
      qError("failed to reallocate memory for output buffer, size:%zu", *outBufSize + leftLen);
×
151
      return terrno;
×
152
    }
153
    (void)memcpy(*outBuf + *outBufSize, inBuf, leftLen);
×
154
    inBuf += leftLen;
×
155
    processedSize += leftLen;
×
156
    *outBufSize += leftLen;
×
157
  }
158

159
  qTrace("[StreamInternal] get result inBufSize:%zu, outBufSize:%zu", inBufSize, *outBufSize);
×
160
  return TSDB_CODE_SUCCESS;
×
161
}
162

163
// Convert result row to buf for rocksdb
164
int32_t putResultRowToBuf(SExprSupp* pSup, const char* inBuf, size_t inBufSize, char** outBuf, size_t* outBufSize) {
×
165
  if (pSup == NULL || inBuf == NULL || outBuf == NULL || outBufSize == NULL) {
×
166
    qError("invalid input parameters, inBuf:%p, pSup:%p, outBufSize:%p, outBuf:%p", inBuf, pSup, outBufSize, outBuf);
×
167
    return TSDB_CODE_INVALID_PARA;
×
168
  }
169

170
  SqlFunctionCtx* pCtx = pSup->pCtx;
×
171
  int32_t*        offset = pSup->rowEntryInfoOffset;
×
172
  SResultRow*     pResultRow = (SResultRow*)inBuf;
×
173
  size_t          rowSize = getResultRowSize(pCtx, pSup->numOfExprs);
×
174

175
  if (rowSize > inBufSize) {
×
176
    qError("invalid input buffer size, rowSize:%zu, inBufSize:%zu", rowSize, inBufSize);
×
177
    return TSDB_CODE_INVALID_PARA;
×
178
  }
179

180
  // calculate the size of output buffer
181
  *outBufSize = rowSize + sizeof(int32_t) * pSup->numOfExprs;
×
182
  if (rowSize < inBufSize) {
×
183
    *outBufSize += inBufSize - rowSize;
×
184
  }
185

186
  *outBuf = taosMemoryMalloc(*outBufSize);
×
187
  if (*outBuf == NULL) {
×
188
    qError("failed to allocate memory for output buffer, size:%zu", *outBufSize);
×
189
    return terrno;
×
190
  }
191

192
  char* pBuf = *outBuf;
×
193
  pResultRow->version = FUNCTION_RESULT_INFO_VERSION;
×
194
  (void)memcpy(pBuf, pResultRow, sizeof(SResultRow));
×
195
  pBuf += sizeof(SResultRow);
×
196
  for (int32_t i = 0; i < pSup->numOfExprs; ++i) {
×
197
    size_t len = sizeof(SResultRowEntryInfo) + pCtx[i].resDataInfo.interBufSize;
×
198
    *(int32_t*)pBuf = (int32_t)len;
×
199
    pBuf += sizeof(int32_t);
×
200
    (void)memcpy(pBuf, getResultEntryInfo(pResultRow, i, offset), len);
×
201
    pBuf += len;
×
202
  }
203

204
  if (rowSize < inBufSize) {
×
205
    // stream stores extra data after result row
206
    size_t leftLen = inBufSize - rowSize;
×
207
    (void)memcpy(pBuf, inBuf + rowSize, leftLen);
×
208
    pBuf += leftLen;
×
209
  }
210

211
  qTrace("[StreamInternal] put result inBufSize:%zu, outBufSize:%zu", inBufSize, *outBufSize);
×
212
  return TSDB_CODE_SUCCESS;
×
213
}
214

215
static void freeEx(void* p) { taosMemoryFree(*(void**)p); }
×
216

217
void cleanupGroupResInfo(SGroupResInfo* pGroupResInfo) {
63,699,773✔
218
  taosMemoryFreeClear(pGroupResInfo->pBuf);
63,699,773✔
219
  if (pGroupResInfo->freeItem) {
63,705,004✔
220
    //    taosArrayDestroy(pGroupResInfo->pRows);
221
    taosArrayDestroyEx(pGroupResInfo->pRows, freeEx);
×
222
    pGroupResInfo->freeItem = false;
×
223
    pGroupResInfo->pRows = NULL;
×
224
  } else {
225
    taosArrayDestroy(pGroupResInfo->pRows);
63,701,060✔
226
    pGroupResInfo->pRows = NULL;
63,702,680✔
227
  }
228
  pGroupResInfo->index = 0;
63,703,284✔
229
  pGroupResInfo->delIndex = 0;
63,703,434✔
230
}
63,698,977✔
231

232
int32_t resultrowComparAsc(const void* p1, const void* p2) {
2,147,483,647✔
233
  SResKeyPos* pp1 = *(SResKeyPos**)p1;
2,147,483,647✔
234
  SResKeyPos* pp2 = *(SResKeyPos**)p2;
2,147,483,647✔
235

236
  if (pp1->groupId == pp2->groupId) {
2,147,483,647✔
237
    int64_t pts1 = *(int64_t*)pp1->key;
2,147,483,647✔
238
    int64_t pts2 = *(int64_t*)pp2->key;
2,147,483,647✔
239

240
    if (pts1 == pts2) {
2,147,483,647✔
241
      return 0;
×
242
    } else {
243
      return pts1 < pts2 ? -1 : 1;
2,147,483,647✔
244
    }
245
  } else {
246
    return pp1->groupId < pp2->groupId ? -1 : 1;
2,147,483,647✔
247
  }
248
}
249

250
static int32_t resultrowComparDesc(const void* p1, const void* p2) { return resultrowComparAsc(p2, p1); }
1,896,596,928✔
251

252
int32_t initGroupedResultInfo(SGroupResInfo* pGroupResInfo, SSHashObj* pHashmap, int32_t order) {
45,366,435✔
253
  int32_t code = TSDB_CODE_SUCCESS;
45,366,435✔
254
  int32_t lino = 0;
45,366,435✔
255
  if (pGroupResInfo->pRows != NULL) {
45,366,435✔
256
    taosArrayDestroy(pGroupResInfo->pRows);
3,285,224✔
257
  }
258
  if (pGroupResInfo->pBuf) {
45,371,545✔
259
    taosMemoryFree(pGroupResInfo->pBuf);
3,285,224✔
260
    pGroupResInfo->pBuf = NULL;
3,285,224✔
261
  }
262

263
  // extract the result rows information from the hash map
264
  int32_t size = tSimpleHashGetSize(pHashmap);
45,353,380✔
265

266
  void* pData = NULL;
45,361,846✔
267
  pGroupResInfo->pRows = taosArrayInit(size, POINTER_BYTES);
45,361,846✔
268
  QUERY_CHECK_NULL(pGroupResInfo->pRows, code, lino, _end, terrno);
45,370,159✔
269

270
  size_t  keyLen = 0;
45,366,737✔
271
  int32_t iter = 0;
45,366,928✔
272
  int64_t bufLen = 0, offset = 0;
45,365,912✔
273

274
  // todo move away and record this during create window
275
  while ((pData = tSimpleHashIterate(pHashmap, pData, &iter)) != NULL) {
2,147,483,647✔
276
    /*void* key = */ (void)tSimpleHashGetKey(pData, &keyLen);
277
    bufLen += keyLen + sizeof(SResultRowPosition);
2,147,483,647✔
278
  }
279

280
  pGroupResInfo->pBuf = taosMemoryMalloc(bufLen);
45,297,499✔
281
  QUERY_CHECK_NULL(pGroupResInfo->pBuf, code, lino, _end, terrno);
45,372,556✔
282

283
  iter = 0;
45,371,245✔
284
  while ((pData = tSimpleHashIterate(pHashmap, pData, &iter)) != NULL) {
2,147,483,647✔
285
    void* key = tSimpleHashGetKey(pData, &keyLen);
2,147,483,647✔
286

287
    SResKeyPos* p = (SResKeyPos*)(pGroupResInfo->pBuf + offset);
2,147,483,647✔
288

289
    p->groupId = *(uint64_t*)key;
2,147,483,647✔
290
    p->pos = *(SResultRowPosition*)pData;
2,147,483,647✔
291
    memcpy(p->key, (char*)key + sizeof(uint64_t), keyLen - sizeof(uint64_t));
2,147,483,647✔
292
    void* tmp = taosArrayPush(pGroupResInfo->pRows, &p);
2,147,483,647✔
293
    QUERY_CHECK_NULL(pGroupResInfo->pBuf, code, lino, _end, terrno);
2,147,483,647✔
294

295
    offset += keyLen + sizeof(struct SResultRowPosition);
2,147,483,647✔
296
  }
297

298
  if (order == TSDB_ORDER_ASC || order == TSDB_ORDER_DESC) {
45,348,061✔
299
    __compar_fn_t fn = (order == TSDB_ORDER_ASC) ? resultrowComparAsc : resultrowComparDesc;
2,857,357✔
300
    size = POINTER_BYTES;
2,857,357✔
301
    taosSort(pGroupResInfo->pRows->pData, taosArrayGetSize(pGroupResInfo->pRows), size, fn);
2,857,357✔
302
  }
303

304
  pGroupResInfo->index = 0;
45,348,542✔
305

306
_end:
45,367,177✔
307
  if (code != TSDB_CODE_SUCCESS) {
45,374,950✔
308
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
309
  }
310
  return code;
45,374,950✔
311
}
312

313
void initMultiResInfoFromArrayList(SGroupResInfo* pGroupResInfo, SArray* pArrayList) {
×
314
  if (pGroupResInfo->pRows != NULL) {
×
315
    taosArrayDestroy(pGroupResInfo->pRows);
×
316
  }
317

318
  pGroupResInfo->freeItem = true;
×
319
  pGroupResInfo->pRows = pArrayList;
×
320
  pGroupResInfo->index = 0;
×
321
  pGroupResInfo->delIndex = 0;
×
322
}
×
323

324
bool hasRemainResults(SGroupResInfo* pGroupResInfo) {
164,711,354✔
325
  if (pGroupResInfo->pRows == NULL) {
164,711,354✔
326
    return false;
×
327
  }
328

329
  return pGroupResInfo->index < taosArrayGetSize(pGroupResInfo->pRows);
164,719,781✔
330
}
331

332
int32_t getNumOfTotalRes(SGroupResInfo* pGroupResInfo) {
89,546,541✔
333
  if (pGroupResInfo->pRows == 0) {
89,546,541✔
334
    return 0;
×
335
  }
336

337
  return (int32_t)taosArrayGetSize(pGroupResInfo->pRows);
89,551,544✔
338
}
339

340
SArray* createSortInfo(SNodeList* pNodeList) {
23,288,009✔
341
  size_t numOfCols = 0;
23,288,009✔
342

343
  if (pNodeList != NULL) {
23,288,009✔
344
    numOfCols = LIST_LENGTH(pNodeList);
23,252,553✔
345
  } else {
346
    numOfCols = 0;
35,604✔
347
  }
348

349
  SArray* pList = taosArrayInit(numOfCols, sizeof(SBlockOrderInfo));
23,290,557✔
350
  if (pList == NULL) {
23,287,820✔
351
    return pList;
×
352
  }
353

354
  for (int32_t i = 0; i < numOfCols; ++i) {
50,689,630✔
355
    SOrderByExprNode* pSortKey = (SOrderByExprNode*)nodesListGetNode(pNodeList, i);
27,399,240✔
356
    if (!pSortKey) {
27,401,592✔
357
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
358
      taosArrayDestroy(pList);
×
359
      pList = NULL;
×
360
      terrno = TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR;
×
361
      break;
×
362
    }
363
    SBlockOrderInfo bi = {0};
27,401,592✔
364
    bi.order = (pSortKey->order == ORDER_ASC) ? TSDB_ORDER_ASC : TSDB_ORDER_DESC;
27,400,121✔
365
    bi.nullFirst = (pSortKey->nullOrder == NULL_ORDER_FIRST);
27,402,275✔
366

367
    if (nodeType(pSortKey->pExpr) != QUERY_NODE_COLUMN) {
27,400,234✔
368
      qError("invalid order by expr type:%d", nodeType(pSortKey->pExpr));
×
369
      taosArrayDestroy(pList);
×
370
      pList = NULL;
×
371
      terrno = TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR;
×
372
      break;
×
373
    }
374
    
375
    SColumnNode* pColNode = (SColumnNode*)pSortKey->pExpr;
27,393,567✔
376
    bi.slotId = pColNode->slotId;
27,396,479✔
377
    void* tmp = taosArrayPush(pList, &bi);
27,401,810✔
378
    if (!tmp) {
27,401,810✔
379
      taosArrayDestroy(pList);
×
380
      pList = NULL;
×
381
      break;
×
382
    }
383
  }
384

385
  return pList;
23,290,390✔
386
}
387

388
SSDataBlock* createDataBlockFromDescNode(void* p) {
337,237,501✔
389
  SDataBlockDescNode* pNode = (SDataBlockDescNode*)p;
337,237,501✔
390
  int32_t      numOfCols = LIST_LENGTH(pNode->pSlots);
337,237,501✔
391
  SSDataBlock* pBlock = NULL;
337,276,358✔
392
  int32_t      code = createDataBlock(&pBlock);
337,279,188✔
393
  if (code) {
337,218,591✔
394
    terrno = code;
×
395
    return NULL;
×
396
  }
397

398
  pBlock->info.id.blockId = pNode->dataBlockId;
337,218,591✔
399
  pBlock->info.type = STREAM_INVALID;
337,217,575✔
400
  pBlock->info.calWin = (STimeWindow){.skey = INT64_MIN, .ekey = INT64_MAX};
337,243,410✔
401
  pBlock->info.watermark = INT64_MIN;
337,225,547✔
402

403
  for (int32_t i = 0; i < numOfCols; ++i) {
1,968,637,812✔
404
    SSlotDescNode* pDescNode = (SSlotDescNode*)nodesListGetNode(pNode->pSlots, i);
1,631,323,766✔
405
    if (!pDescNode) {
1,631,166,784✔
406
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(code));
×
407
      blockDataDestroy(pBlock);
×
408
      pBlock = NULL;
×
409
      terrno = TSDB_CODE_INVALID_PARA;
×
410
      break;
×
411
    }
412
    SColumnInfoData idata =
1,631,110,566✔
413
        createColumnInfoData(pDescNode->dataType.type, pDescNode->dataType.bytes, pDescNode->slotId);
1,631,198,533✔
414
    idata.info.scale = pDescNode->dataType.scale;
1,631,453,002✔
415
    idata.info.precision = pDescNode->dataType.precision;
1,631,471,902✔
416
    idata.info.noData = pDescNode->reserve;
1,631,513,507✔
417

418
    code = blockDataAppendColInfo(pBlock, &idata);
1,631,550,231✔
419
    if (code != TSDB_CODE_SUCCESS) {
1,631,455,315✔
420
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(code));
208,110✔
421
      blockDataDestroy(pBlock);
208,110✔
422
      pBlock = NULL;
×
423
      terrno = code;
×
424
      break;
×
425
    }
426
  }
427

428
  return pBlock;
337,339,089✔
429
}
430

431
int32_t prepareDataBlockBuf(SSDataBlock* pDataBlock, SColMatchInfo* pMatchInfo) {
113,539,114✔
432
  SDataBlockInfo* pBlockInfo = &pDataBlock->info;
113,539,114✔
433

434
  for (int32_t i = 0; i < taosArrayGetSize(pMatchInfo->pList); ++i) {
693,591,110✔
435
    SColMatchItem* pItem = taosArrayGet(pMatchInfo->pList, i);
587,481,306✔
436
    if (!pItem) {
587,430,767✔
437
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
438
      return terrno;
×
439
    }
440

441
    if (pItem->isPk) {
587,430,767✔
442
      SColumnInfoData* pInfoData = taosArrayGet(pDataBlock->pDataBlock, pItem->dstSlotId);
7,431,649✔
443
      if (!pInfoData) {
7,356,200✔
444
        qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
445
        return terrno;
×
446
      }
447
      pBlockInfo->pks[0].type = pInfoData->info.type;
7,356,200✔
448
      pBlockInfo->pks[1].type = pInfoData->info.type;
7,375,610✔
449

450
      // allocate enough buffer size, which is pInfoData->info.bytes
451
      if (IS_VAR_DATA_TYPE(pItem->dataType.type)) {
7,370,672✔
452
        pBlockInfo->pks[0].pData = taosMemoryCalloc(1, pInfoData->info.bytes);
2,464,073✔
453
        if (pBlockInfo->pks[0].pData == NULL) {
2,461,499✔
454
          return terrno;
×
455
        }
456

457
        pBlockInfo->pks[1].pData = taosMemoryCalloc(1, pInfoData->info.bytes);
2,464,733✔
458
        if (pBlockInfo->pks[1].pData == NULL) {
2,462,009✔
459
          taosMemoryFreeClear(pBlockInfo->pks[0].pData);
×
460
          return terrno;
×
461
        }
462

463
        pBlockInfo->pks[0].nData = pInfoData->info.bytes;
2,463,371✔
464
        pBlockInfo->pks[1].nData = pInfoData->info.bytes;
2,467,628✔
465
      }
466

467
      break;
7,371,884✔
468
    }
469
  }
470

471
  return TSDB_CODE_SUCCESS;
113,515,696✔
472
}
473

474
EDealRes doTranslateTagExpr(SNode** pNode, void* pContext) {
391,300✔
475
  STransTagExprCtx* pCtx = pContext;
391,300✔
476
  SMetaReader*      mr = pCtx->pReader;
391,300✔
477
  bool              isTagCol = false, isTbname = false;
391,300✔
478
  if (nodeType(*pNode) == QUERY_NODE_COLUMN) {
391,300✔
479
    SColumnNode* pCol = (SColumnNode*)*pNode;
111,800✔
480
    if (pCol->colType == COLUMN_TYPE_TBNAME)
111,800✔
481
      isTbname = true;
×
482
    else
483
      isTagCol = true;
111,800✔
484
  } else if (nodeType(*pNode) == QUERY_NODE_FUNCTION) {
279,500✔
485
    SFunctionNode* pFunc = (SFunctionNode*)*pNode;
×
486
    if (pFunc->funcType == FUNCTION_TYPE_TBNAME) isTbname = true;
×
487
  }
488
  if (isTagCol) {
391,300✔
489
    SColumnNode* pSColumnNode = *(SColumnNode**)pNode;
111,800✔
490

491
    SValueNode* res = NULL;
111,800✔
492
    pCtx->code = nodesMakeNode(QUERY_NODE_VALUE, (SNode**)&res);
111,800✔
493
    if (NULL == res) {
111,800✔
494
      return DEAL_RES_ERROR;
×
495
    }
496

497
    res->translate = true;
111,800✔
498
    res->node.resType = pSColumnNode->node.resType;
111,800✔
499

500
    STagVal tagVal = {0};
111,800✔
501
    tagVal.cid = pSColumnNode->colId;
111,800✔
502
    const char* p = mr->pAPI->extractTagVal(mr->me.ctbEntry.pTags, pSColumnNode->node.resType.type, &tagVal);
111,800✔
503
    if (p == NULL) {
111,800✔
504
      res->node.resType.type = TSDB_DATA_TYPE_NULL;
×
505
    } else if (pSColumnNode->node.resType.type == TSDB_DATA_TYPE_JSON) {
111,800✔
506
      int32_t len = ((const STag*)p)->len;
×
507
      res->datum.p = taosMemoryCalloc(len + 1, 1);
×
508
      if (NULL == res->datum.p) {
×
509
        return DEAL_RES_ERROR;
×
510
      }
511
      memcpy(res->datum.p, p, len);
×
512
    } else if (IS_VAR_DATA_TYPE(pSColumnNode->node.resType.type)) {
111,800✔
513
      if (IS_STR_DATA_BLOB(pSColumnNode->node.resType.type)) {
111,800✔
514
        return TSDB_CODE_BLOB_NOT_SUPPORT_TAG;
×
515
      }
516

517
      res->datum.p = taosMemoryCalloc(tagVal.nData + VARSTR_HEADER_SIZE + 1, 1);
111,800✔
518
      if (NULL == res->datum.p) {
111,800✔
519
        return DEAL_RES_ERROR;
×
520
      }
521

522
      if (IS_STR_DATA_BLOB(pSColumnNode->node.resType.type)) {
111,800✔
523
        memcpy(blobDataVal(res->datum.p), tagVal.pData, tagVal.nData);
×
524
        blobDataSetLen(res->datum.p, tagVal.nData);
×
525
      } else {
526
        memcpy(varDataVal(res->datum.p), tagVal.pData, tagVal.nData);
111,800✔
527
        varDataSetLen(res->datum.p, tagVal.nData);
111,800✔
528
      }
529
    } else {
530
      int32_t code = nodesSetValueNodeValue(res, &(tagVal.i64));
×
531
      if (code != TSDB_CODE_SUCCESS) {
×
532
        return DEAL_RES_ERROR;
×
533
      }
534
    }
535
    nodesDestroyNode(*pNode);
111,800✔
536
    *pNode = (SNode*)res;
111,800✔
537
  } else if (isTbname) {
279,500✔
538
    SValueNode* res = NULL;
×
539
    pCtx->code = nodesMakeNode(QUERY_NODE_VALUE, (SNode**)&res);
×
540
    if (NULL == res) {
×
541
      return DEAL_RES_ERROR;
×
542
    }
543

544
    res->translate = true;
×
545
    res->node.resType = ((SExprNode*)(*pNode))->resType;
×
546

547
    int32_t len = strlen(mr->me.name);
×
548
    res->datum.p = taosMemoryCalloc(len + VARSTR_HEADER_SIZE + 1, 1);
×
549
    if (NULL == res->datum.p) {
×
550
      return DEAL_RES_ERROR;
×
551
    }
552
    memcpy(varDataVal(res->datum.p), mr->me.name, len);
×
553
    varDataSetLen(res->datum.p, len);
×
554
    nodesDestroyNode(*pNode);
×
555
    *pNode = (SNode*)res;
×
556
  }
557

558
  return DEAL_RES_CONTINUE;
391,300✔
559
}
560

561
int32_t isQualifiedTable(STableKeyInfo* info, SNode* pTagCond, void* metaHandle, bool* pQualified, SStorageAPI* pAPI) {
55,900✔
562
  int32_t     code = TSDB_CODE_SUCCESS;
55,900✔
563
  SMetaReader mr = {0};
55,900✔
564

565
  pAPI->metaReaderFn.initReader(&mr, metaHandle, META_READER_LOCK, &pAPI->metaFn);
55,900✔
566
  code = pAPI->metaReaderFn.getEntryGetUidCache(&mr, info->uid);
55,900✔
567
  if (TSDB_CODE_SUCCESS != code) {
55,900✔
568
    pAPI->metaReaderFn.clearReader(&mr);
×
569
    *pQualified = false;
×
570

571
    return TSDB_CODE_SUCCESS;
×
572
  }
573

574
  SNode* pTagCondTmp = NULL;
55,900✔
575
  code = nodesCloneNode(pTagCond, &pTagCondTmp);
55,900✔
576
  if (TSDB_CODE_SUCCESS != code) {
55,900✔
577
    *pQualified = false;
×
578
    pAPI->metaReaderFn.clearReader(&mr);
×
579
    return code;
×
580
  }
581
  STransTagExprCtx ctx = {.code = 0, .pReader = &mr};
55,900✔
582
  nodesRewriteExprPostOrder(&pTagCondTmp, doTranslateTagExpr, &ctx);
55,900✔
583
  pAPI->metaReaderFn.clearReader(&mr);
55,900✔
584
  if (TSDB_CODE_SUCCESS != ctx.code) {
55,900✔
585
    *pQualified = false;
×
586
    terrno = code;
×
587
    return code;
×
588
  }
589

590
  SNode* pNew = NULL;
55,900✔
591
  code = scalarCalculateConstants(pTagCondTmp, &pNew);
55,900✔
592
  if (TSDB_CODE_SUCCESS != code) {
55,900✔
593
    terrno = code;
×
594
    nodesDestroyNode(pTagCondTmp);
×
595
    *pQualified = false;
×
596

597
    return code;
×
598
  }
599

600
  SValueNode* pValue = (SValueNode*)pNew;
55,900✔
601
  *pQualified = pValue->datum.b;
55,900✔
602

603
  nodesDestroyNode(pNew);
55,900✔
604
  return TSDB_CODE_SUCCESS;
55,900✔
605
}
606

607
static EDealRes getColumn(SNode** pNode, void* pContext) {
46,031,917✔
608
  tagFilterAssist* pData = (tagFilterAssist*)pContext;
46,031,917✔
609
  SColumnNode*     pSColumnNode = NULL;
46,031,917✔
610
  if (QUERY_NODE_COLUMN == nodeType((*pNode))) {
46,032,454✔
611
    pSColumnNode = *(SColumnNode**)pNode;
15,085,468✔
612
  } else if (QUERY_NODE_FUNCTION == nodeType((*pNode))) {
30,948,867✔
613
    SFunctionNode* pFuncNode = *(SFunctionNode**)(pNode);
670,344✔
614
    if (pFuncNode->funcType == FUNCTION_TYPE_TBNAME) {
670,344✔
615
      pData->code = nodesMakeNode(QUERY_NODE_COLUMN, (SNode**)&pSColumnNode);
623,153✔
616
      if (NULL == pSColumnNode) {
623,153✔
617
        return DEAL_RES_ERROR;
×
618
      }
619
      pSColumnNode->colId = -1;
623,153✔
620
      pSColumnNode->colType = COLUMN_TYPE_TBNAME;
623,153✔
621
      pSColumnNode->node.resType.type = TSDB_DATA_TYPE_VARCHAR;
623,153✔
622
      pSColumnNode->node.resType.bytes = TSDB_TABLE_FNAME_LEN - 1 + VARSTR_HEADER_SIZE;
623,153✔
623
      nodesDestroyNode(*pNode);
623,153✔
624
      *pNode = (SNode*)pSColumnNode;
623,153✔
625
    } else {
626
      return DEAL_RES_CONTINUE;
47,191✔
627
    }
628
  } else {
629
    return DEAL_RES_CONTINUE;
30,278,434✔
630
  }
631

632
  void* data = taosHashGet(pData->colHash, &pSColumnNode->colId, sizeof(pSColumnNode->colId));
15,708,504✔
633
  if (!data) {
15,707,701✔
634
    int32_t tempRes =
635
        taosHashPut(pData->colHash, &pSColumnNode->colId, sizeof(pSColumnNode->colId), pNode, sizeof((*pNode)));
14,388,361✔
636
    if (tempRes != TSDB_CODE_SUCCESS && tempRes != TSDB_CODE_DUP_KEY) {
14,388,761✔
637
      return DEAL_RES_ERROR;
×
638
    }
639
    pSColumnNode->slotId = pData->index++;
14,388,761✔
640
    SColumnInfo cInfo = {.colId = pSColumnNode->colId,
14,387,811✔
641
                         .type = pSColumnNode->node.resType.type,
14,388,074✔
642
                         .bytes = pSColumnNode->node.resType.bytes,
14,389,627✔
643
                         .pk = pSColumnNode->isPk};
14,388,151✔
644
#if TAG_FILTER_DEBUG
645
    qDebug("tagfilter build column info, slotId:%d, colId:%d, type:%d", pSColumnNode->slotId, cInfo.colId, cInfo.type);
646
#endif
647
    void* tmp = taosArrayPush(pData->cInfoList, &cInfo);
14,386,929✔
648
    if (!tmp) {
14,389,824✔
649
      return DEAL_RES_ERROR;
×
650
    }
651
  } else {
652
    SColumnNode* col = *(SColumnNode**)data;
1,319,340✔
653
    pSColumnNode->slotId = col->slotId;
1,319,340✔
654
  }
655

656
  return DEAL_RES_CONTINUE;
15,706,417✔
657
}
658

659
static int32_t createResultData(SDataType* pType, int32_t numOfRows, SScalarParam* pParam) {
13,248,864✔
660
  SColumnInfoData* pColumnData = taosMemoryCalloc(1, sizeof(SColumnInfoData));
13,248,864✔
661
  if (pColumnData == NULL) {
13,251,323✔
662
    return terrno;
×
663
  }
664

665
  pColumnData->info.type = pType->type;
13,251,323✔
666
  pColumnData->info.bytes = pType->bytes;
13,248,985✔
667
  pColumnData->info.scale = pType->scale;
13,249,615✔
668
  pColumnData->info.precision = pType->precision;
13,248,994✔
669

670
  int32_t code = colInfoDataEnsureCapacity(pColumnData, numOfRows, true);
13,248,734✔
671
  if (code != TSDB_CODE_SUCCESS) {
13,247,152✔
672
    terrno = code;
×
673
    releaseColInfoData(pColumnData);
×
674
    return terrno;
×
675
  }
676

677
  pParam->columnData = pColumnData;
13,247,152✔
678
  pParam->colAlloced = true;
13,246,788✔
679
  return TSDB_CODE_SUCCESS;
13,248,925✔
680
}
681

682
static void releaseColInfoData(void* pCol) {
2,158,154✔
683
  if (pCol) {
2,158,154✔
684
    SColumnInfoData* col = (SColumnInfoData*)pCol;
2,158,154✔
685
    colDataDestroy(col);
2,158,154✔
686
    taosMemoryFree(col);
2,158,154✔
687
  }
688
}
2,158,154✔
689

690
void freeItem(void* p) {
186,387,554✔
691
  STUidTagInfo* pInfo = p;
186,387,554✔
692
  if (pInfo->pTagVal != NULL) {
186,387,554✔
693
    taosMemoryFree(pInfo->pTagVal);
185,965,786✔
694
  }
695
}
186,387,415✔
696

697
typedef struct {
698
  col_id_t  colId;
699
  SNode*    pValueNode;
700
  int32_t   bytes;  // length defined in schema
701
} STagDataEntry;
702

703
static int compareTagDataEntry(const void* a, const void* b) {
40,040✔
704
  STagDataEntry* p1 = (STagDataEntry*)a;
40,040✔
705
  STagDataEntry* p2 = (STagDataEntry*)b;
40,040✔
706
  return compareInt16Val(&p1->colId, &p2->colId);
40,040✔
707
}
708

709
static int32_t buildTagDataEntryKey(SArray* pIdWithValue, char** keyBuf, int32_t keyLen) {
20,020✔
710
  *keyBuf = (char*)taosMemoryCalloc(1, keyLen);
20,020✔
711
  if (NULL == *keyBuf) {
20,020✔
712
    qError(
×
713
      "failed to allocate memory for tag filter optimization key, size:%d",
714
      keyLen);
715
    return terrno;
×
716
  }
717
  char* pStart = *keyBuf;
20,020✔
718
  for (int32_t i = 0; i < taosArrayGetSize(pIdWithValue); ++i) {
60,060✔
719
    STagDataEntry* entry      = (STagDataEntry*)taosArrayGet(pIdWithValue, i);
40,040✔
720
    SValueNode*    pValueNode = (SValueNode*)entry->pValueNode;
40,040✔
721
    // num type may have different bytes length, use the smaller one
722
    int32_t        bytes = TMIN(entry->bytes, pValueNode->node.resType.bytes);
40,040✔
723

724
    (void)memcpy(pStart, &entry->colId, sizeof(col_id_t));
40,040✔
725
    pStart += sizeof(col_id_t);
40,040✔
726

727
    if (!pValueNode->isNull) {
40,040✔
728
      switch (pValueNode->node.resType.type) {
34,580✔
729
        case TSDB_DATA_TYPE_BOOL:
4,004✔
730
          (void)memcpy(
4,004✔
731
            pStart, &pValueNode->datum.b, bytes);
4,004✔
732
          pStart += bytes;
4,004✔
733
          break;
4,004✔
734
        case TSDB_DATA_TYPE_TINYINT:
16,380✔
735
        case TSDB_DATA_TYPE_SMALLINT:
736
        case TSDB_DATA_TYPE_INT:
737
        case TSDB_DATA_TYPE_BIGINT:
738
        case TSDB_DATA_TYPE_TIMESTAMP:
739
          (void)memcpy(
16,380✔
740
            pStart, &pValueNode->datum.i, bytes);
16,380✔
741
          pStart += bytes;
16,380✔
742
          break;
16,380✔
743
        case TSDB_DATA_TYPE_UTINYINT:
×
744
        case TSDB_DATA_TYPE_USMALLINT:
745
        case TSDB_DATA_TYPE_UINT:
746
        case TSDB_DATA_TYPE_UBIGINT:
747
          (void)memcpy(
×
748
            pStart, &pValueNode->datum.u, bytes);
×
749
          pStart += bytes;
×
750
          break;
×
751
        case TSDB_DATA_TYPE_FLOAT:
4,732✔
752
        case TSDB_DATA_TYPE_DOUBLE:
753
          (void)memcpy(
4,732✔
754
            pStart, &pValueNode->datum.d, bytes);
4,732✔
755
          pStart += bytes;
4,732✔
756
          break;
4,732✔
757
        case TSDB_DATA_TYPE_VARCHAR:
9,464✔
758
        case TSDB_DATA_TYPE_VARBINARY:
759
        case TSDB_DATA_TYPE_NCHAR:
760
          (void)memcpy(pStart,
9,464✔
761
            varDataVal(pValueNode->datum.p), varDataLen(pValueNode->datum.p));
9,464✔
762
          pStart += varDataLen(pValueNode->datum.p);
9,464✔
763
          break;
9,464✔
764
        default:
×
765
          qError("unsupported tag data type %d in tag filter optimization",
×
766
            pValueNode->node.resType.type);
767
          return TSDB_CODE_STREAM_INTERNAL_ERROR;
×
768
      }
769
    }
770
  }
771

772
  return TSDB_CODE_SUCCESS;
20,020✔
773
}
774

775
static void extractTagDataEntry(
40,040✔
776
  SOperatorNode* pOpNode, SArray* pIdWithValue) {
777
  SNode* pLeft = pOpNode->pLeft;
40,040✔
778
  SNode* pRight = pOpNode->pRight;
40,040✔
779
  SColumnNode* pColNode = nodeType(pLeft) == QUERY_NODE_COLUMN ?
40,040✔
780
    (SColumnNode*)pLeft : (SColumnNode*)pRight;
40,040✔
781
  SValueNode* pValueNode = nodeType(pLeft) == QUERY_NODE_VALUE ?
40,040✔
782
    (SValueNode*)pLeft : (SValueNode*)pRight;
40,040✔
783

784
  STagDataEntry entry = {0};
40,040✔
785
  entry.colId = pColNode->colId;
40,040✔
786
  entry.pValueNode = (SNode*)pValueNode;
40,040✔
787
  entry.bytes = pColNode->node.resType.bytes;
40,040✔
788
  void* _tmp = taosArrayPush(pIdWithValue, &entry);
40,040✔
789
}
40,040✔
790

791
static int32_t extractTagFilterTagDataEntries(
20,020✔
792
  const SNode* pTagCond, SArray* pIdWithVal) {
793
  if (NULL == pTagCond || NULL == pIdWithVal ||
20,020✔
794
    (nodeType(pTagCond) != QUERY_NODE_OPERATOR &&
20,020✔
795
      nodeType(pTagCond) != QUERY_NODE_LOGIC_CONDITION)) {
20,020✔
796
    qError("invalid parameter to extract tag filter symbol");
×
797
    return TSDB_CODE_STREAM_INTERNAL_ERROR;
×
798
  }
799

800
  if (nodeType(pTagCond) == QUERY_NODE_OPERATOR) {
20,020✔
801
    extractTagDataEntry((SOperatorNode*)pTagCond, pIdWithVal);
×
802
  } else if (nodeType(pTagCond) == QUERY_NODE_LOGIC_CONDITION) {
20,020✔
803
    SNode* pChild = NULL;
20,020✔
804
    FOREACH(pChild, ((SLogicConditionNode*)pTagCond)->pParameterList) {
60,060✔
805
      extractTagDataEntry((SOperatorNode*)pChild, pIdWithVal);
40,040✔
806
    }
807
  }
808

809
  taosArraySort(pIdWithVal, compareTagDataEntry);
20,020✔
810

811
  return TSDB_CODE_SUCCESS;
20,020✔
812
}
813

814
static int32_t genStableTagFilterDigest(const SNode* pTagCond, T_MD5_CTX* pContext) {
20,020✔
815
  if (pTagCond == NULL) {
20,020✔
816
    return TSDB_CODE_SUCCESS;
×
817
  }
818

819
  char*   payload = NULL;
20,020✔
820
  int32_t len = 0;
20,020✔
821
  int32_t code = TSDB_CODE_SUCCESS;
20,020✔
822
  int32_t lino = 0;
20,020✔
823

824
  SArray* pIdWithVal = taosArrayInit(TARRAY_MIN_SIZE, sizeof(STagDataEntry));
20,020✔
825
  code = extractTagFilterTagDataEntries(pTagCond, pIdWithVal);
20,020✔
826
  QUERY_CHECK_CODE(code, lino, _end);
20,020✔
827
  for (int32_t i = 0; i < taosArrayGetSize(pIdWithVal); ++i) {
60,060✔
828
    STagDataEntry* pEntry = taosArrayGet(pIdWithVal, i);
40,040✔
829
    len += sizeof(col_id_t) + pEntry->bytes;
40,040✔
830
  }
831
  code = buildTagDataEntryKey(pIdWithVal, &payload, len);
20,020✔
832
  QUERY_CHECK_CODE(code, lino, _end);
20,020✔
833

834
  tMD5Init(pContext);
20,020✔
835
  tMD5Update(pContext, (uint8_t*)payload, (uint32_t)len);
20,020✔
836
  tMD5Final(pContext);
20,020✔
837

838
_end:
20,020✔
839
  if (TSDB_CODE_SUCCESS != code) {
20,020✔
840
    qError("%s failed at line %d since %s",
×
841
      __func__, __LINE__, tstrerror(code));
842
  }
843
  taosArrayDestroy(pIdWithVal);
20,020✔
844
  taosMemoryFree(payload);
20,020✔
845
  return code;
20,020✔
846
}
847

848
static int32_t genTagFilterDigest(const SNode* pTagCond, T_MD5_CTX* pContext) {
64,723✔
849
  if (pTagCond == NULL) {
64,723✔
850
    return TSDB_CODE_SUCCESS;
61,446✔
851
  }
852

853
  char*   payload = NULL;
3,277✔
854
  int32_t len = 0;
3,277✔
855
  int32_t code = nodesNodeToMsg(pTagCond, &payload, &len);
3,277✔
856
  if (code != TSDB_CODE_SUCCESS) {
3,277✔
857
    qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(code));
×
858
    return code;
×
859
  }
860

861
  tMD5Init(pContext);
3,277✔
862
  tMD5Update(pContext, (uint8_t*)payload, (uint32_t)len);
3,277✔
863
  tMD5Final(pContext);
3,277✔
864

865
  // void* tmp = NULL;
866
  // uint32_t size = 0;
867
  // (void)taosAscii2Hex((const char*)pContext->digest, 16, &tmp, &size);
868
  // qInfo("tag filter digest payload: %s", tmp);
869
  // taosMemoryFree(tmp);
870

871
  taosMemoryFree(payload);
3,277✔
872
  return TSDB_CODE_SUCCESS;
3,277✔
873
}
874

875
static int32_t genTbGroupDigest(const SNode* pGroup, uint8_t* filterDigest, T_MD5_CTX* pContext) {
×
876
  int32_t code = TSDB_CODE_SUCCESS;
×
877
  int32_t lino = 0;
×
878
  char*   payload = NULL;
×
879
  int32_t len = 0;
×
880
  code = nodesNodeToMsg(pGroup, &payload, &len);
×
881
  QUERY_CHECK_CODE(code, lino, _end);
×
882

883
  if (filterDigest[0]) {
×
884
    payload = taosMemoryRealloc(payload, len + tListLen(pContext->digest));
×
885
    QUERY_CHECK_NULL(payload, code, lino, _end, terrno);
×
886
    memcpy(payload + len, filterDigest + 1, tListLen(pContext->digest));
×
887
    len += tListLen(pContext->digest);
×
888
  }
889

890
  tMD5Init(pContext);
×
891
  tMD5Update(pContext, (uint8_t*)payload, (uint32_t)len);
×
892
  tMD5Final(pContext);
×
893

894
_end:
×
895
  taosMemoryFree(payload);
×
896
  if (code != TSDB_CODE_SUCCESS) {
×
897
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
898
  }
899
  return code;
×
900
}
901

902
int32_t qGetColumnsFromNodeList(void* data, bool isList, SArray** pColList) {
13,064,003✔
903
  int32_t code = TSDB_CODE_SUCCESS;
13,064,003✔
904
  tagFilterAssist ctx = {0};
13,064,003✔
905
  ctx.colHash = taosHashInit(4, taosGetDefaultHashFunction(TSDB_DATA_TYPE_SMALLINT), false, HASH_NO_LOCK);
13,064,758✔
906
  if (ctx.colHash == NULL) {
13,063,711✔
907
    code = terrno;
×
908
    goto end;
×
909
  }
910

911
  ctx.index = 0;
13,063,711✔
912
  ctx.cInfoList = taosArrayInit(4, sizeof(SColumnInfo));
13,063,711✔
913
  if (ctx.cInfoList == NULL) {
13,064,558✔
914
    code = terrno;
519✔
915
    goto end;
×
916
  }
917

918
  if (isList) {
13,064,039✔
919
    SNode* pNode = NULL;
1,970,511✔
920
    FOREACH(pNode, (SNodeList*)data) {
4,128,329✔
921
      nodesRewriteExprPostOrder(&pNode, getColumn, (void*)&ctx);
2,156,866✔
922
      if (TSDB_CODE_SUCCESS != ctx.code) {
2,157,050✔
923
        code = ctx.code;
×
924
        goto end;
×
925
      }
926
      REPLACE_NODE(pNode);
2,157,050✔
927
    }
928
  } else {
929
    SNode* pNode = (SNode*)data;
11,093,528✔
930
    nodesRewriteExprPostOrder(&pNode, getColumn, (void*)&ctx);
11,093,528✔
931
    if (TSDB_CODE_SUCCESS != ctx.code) {
11,093,448✔
932
      code = ctx.code;
×
933
      goto end;
×
934
    }
935
  }
936
  
937
  if (pColList != NULL) *pColList = ctx.cInfoList;
13,062,528✔
938
  ctx.cInfoList = NULL;
13,064,006✔
939

940
end:
13,065,432✔
941
  taosHashCleanup(ctx.colHash);
13,063,776✔
942
  taosArrayDestroy(ctx.cInfoList);
13,062,860✔
943
  return code;
13,063,498✔
944
}
945

946
static int32_t buildGroupInfo(SColumnInfoData* pValue, int32_t i, SArray* gInfo) {
713,581✔
947
  int32_t code = TSDB_CODE_SUCCESS;
713,581✔
948
  SStreamGroupValue* v = taosArrayReserve(gInfo, 1);
713,581✔
949
  if (v == NULL) {
714,320✔
950
    code = terrno;
×
951
    goto end;
×
952
  }
953
  if (colDataIsNull_s(pValue, i)) {
1,427,537✔
954
    v->isNull = true;
14,560✔
955
  } else {
956
    v->isNull = false;
698,657✔
957
    char* data = colDataGetData(pValue, i);
699,021✔
958
    if (pValue->info.type == TSDB_DATA_TYPE_JSON) {
699,021✔
959
      if (tTagIsJson(data)) {
×
960
        code = TSDB_CODE_QRY_JSON_IN_GROUP_ERROR;
×
961
        goto end;
×
962
      }
963
      if (tTagIsJsonNull(data)) {
×
964
        v->isNull = true;
×
965
        goto end;
×
966
      }
967
      int32_t len = getJsonValueLen(data);
×
968
      v->data.type = pValue->info.type;
×
969
      v->data.nData = len;
×
970
      v->data.pData = taosMemoryCalloc(1, len + 1);
×
971
      if (v->data.pData == NULL) {
×
972
        code = terrno;
×
973
        goto end;
×
974
      }
975
      memcpy(v->data.pData, data, len);
×
976
      qDebug("buildGroupInfo:%d add json data len:%d, data:%s", i, len, (char*)v->data.pData);
×
977
    } else if (IS_VAR_DATA_TYPE(pValue->info.type)) {
698,293✔
978
      if (varDataTLen(data) > pValue->info.bytes) {
449,187✔
979
        code = TSDB_CODE_TDB_INVALID_TABLE_SCHEMA_VER;
×
980
        goto end;
×
981
      }
982
      v->data.type = pValue->info.type;
450,654✔
983
      v->data.nData = varDataLen(data);
450,654✔
984
      v->data.pData = taosMemoryCalloc(1, varDataLen(data) + 1);
450,654✔
985
      if (v->data.pData == NULL) {
450,654✔
986
        code = terrno;
×
987
        goto end;
×
988
      }
989
      memcpy(v->data.pData, varDataVal(data), varDataLen(data));
450,654✔
990
      qDebug("buildGroupInfo:%d add var data type:%d, len:%d, data:%s", i, pValue->info.type, varDataLen(data), (char*)v->data.pData);
450,654✔
991
    } else if (pValue->info.type == TSDB_DATA_TYPE_DECIMAL) {  // reader todo decimal
248,731✔
992
      v->data.type = pValue->info.type;
×
993
      v->data.nData = pValue->info.bytes;
×
994
      v->data.pData = taosMemoryCalloc(1, pValue->info.bytes);
×
995
      if (v->data.pData == NULL) {
×
996
        code = terrno;
×
997
        goto end;
×
998
      }
999
      memcpy(&v->data.pData, data, pValue->info.bytes);
×
1000
      qDebug("buildGroupInfo:%d add data type:%d, data:%"PRId64, i, pValue->info.type, v->data.val);
×
1001
    } else {  // reader todo decimal
1002
      v->data.type = pValue->info.type;
248,003✔
1003
      memcpy(&v->data.val, data, pValue->info.bytes);
248,742✔
1004
      qDebug("buildGroupInfo:%d add data type:%d, data:%"PRId64, i, pValue->info.type, v->data.val);
248,003✔
1005
    }
1006
  }
1007
end:
78,981✔
1008
  if (code != TSDB_CODE_SUCCESS) {
713,913✔
1009
    qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(code));
×
1010
    v->isNull = true;
×
1011
  }
1012
  return code;
713,913✔
1013
}
1014

1015
static void getColInfoResultForGroupbyForStream(void* pVnode, SNodeList* group, STableListInfo* pTableListInfo,
106,024✔
1016
                                   SStorageAPI* pAPI, SHashObj* groupIdMap) {
1017
  int32_t      code = TSDB_CODE_SUCCESS;
106,024✔
1018
  int32_t      lino = 0;
106,024✔
1019
  SArray*      pBlockList = NULL;
106,024✔
1020
  SSDataBlock* pResBlock = NULL;
106,024✔
1021
  SArray*      groupData = NULL;
106,404✔
1022
  SArray*      pUidTagList = NULL;
106,404✔
1023
  SArray*      gInfo = NULL;
106,404✔
1024
  int32_t      tbNameIndex = 0;
106,404✔
1025

1026
  int32_t rows = taosArrayGetSize(pTableListInfo->pTableList);
106,404✔
1027
  if (rows == 0) {
106,404✔
1028
    return;
×
1029
  }
1030

1031
  pUidTagList = taosArrayInit(8, sizeof(STUidTagInfo));
106,404✔
1032
  QUERY_CHECK_NULL(pUidTagList, code, lino, end, terrno);
106,404✔
1033

1034
  for (int32_t i = 0; i < rows; ++i) {
448,062✔
1035
    STableKeyInfo* pkeyInfo = taosArrayGet(pTableListInfo->pTableList, i);
341,658✔
1036
    QUERY_CHECK_NULL(pkeyInfo, code, lino, end, terrno);
341,658✔
1037
    STUidTagInfo info = {.uid = pkeyInfo->uid};
341,658✔
1038
    void*        tmp = taosArrayPush(pUidTagList, &info);
341,658✔
1039
    QUERY_CHECK_NULL(tmp, code, lino, end, terrno);
341,658✔
1040
  }
1041
 
1042
  if (taosArrayGetSize(pUidTagList) > 0) {
106,404✔
1043
    code = pAPI->metaFn.getTableTagsByUid(pVnode, pTableListInfo->idInfo.suid, pUidTagList);
106,404✔
1044
  } else {
1045
    code = pAPI->metaFn.getTableTags(pVnode, pTableListInfo->idInfo.suid, pUidTagList);
×
1046
  }
1047
  if (code != TSDB_CODE_SUCCESS) {
106,404✔
1048
    goto end;
×
1049
  }
1050

1051
  SArray* pColList = NULL;
106,404✔
1052
  code = qGetColumnsFromNodeList(group, true, &pColList);
106,404✔
1053
  if (code != TSDB_CODE_SUCCESS) {
106,404✔
1054
    goto end;
×
1055
  }
1056

1057
  for (int32_t i = 0; i < taosArrayGetSize(pColList); ++i) {
278,423✔
1058
    SColumnInfo* tmp = (SColumnInfo*)taosArrayGet(pColList, i);
172,019✔
1059
    if (tmp != NULL && tmp->colId == -1) {
172,019✔
1060
      tbNameIndex = i;
106,404✔
1061
    }
1062
  }
1063
  
1064
  int32_t numOfTables = taosArrayGetSize(pUidTagList);
106,404✔
1065
  pResBlock = createTagValBlockForFilter(pColList, numOfTables, pUidTagList, pVnode, pAPI);
106,404✔
1066
  taosArrayDestroy(pColList);
106,404✔
1067
  if (pResBlock == NULL) {
106,404✔
1068
    code = terrno;
×
1069
    goto end;
×
1070
  }
1071

1072
  pBlockList = taosArrayInit(2, POINTER_BYTES);
106,404✔
1073
  QUERY_CHECK_NULL(pBlockList, code, lino, end, terrno);
106,404✔
1074

1075
  void* tmp = taosArrayPush(pBlockList, &pResBlock);
106,404✔
1076
  QUERY_CHECK_NULL(tmp, code, lino, end, terrno);
106,404✔
1077

1078
  groupData = taosArrayInit(2, POINTER_BYTES);
106,404✔
1079
  QUERY_CHECK_NULL(groupData, code, lino, end, terrno);
106,404✔
1080

1081
  SNode* pNode = NULL;
106,404✔
1082
  FOREACH(pNode, group) {
278,423✔
1083
    SScalarParam output = {0};
172,019✔
1084

1085
    switch (nodeType(pNode)) {
172,019✔
1086
      case QUERY_NODE_VALUE:
×
1087
        break;
×
1088
      case QUERY_NODE_COLUMN:
172,019✔
1089
      case QUERY_NODE_OPERATOR:
1090
      case QUERY_NODE_FUNCTION: {
1091
        SExprNode* expNode = (SExprNode*)pNode;
172,019✔
1092
        code = createResultData(&expNode->resType, rows, &output);
172,019✔
1093
        if (code != TSDB_CODE_SUCCESS) {
172,019✔
1094
          goto end;
×
1095
        }
1096
        break;
172,019✔
1097
      }
1098

1099
      default:
×
1100
        code = TSDB_CODE_OPS_NOT_SUPPORT;
×
1101
        goto end;
×
1102
    }
1103

1104
    if (nodeType(pNode) == QUERY_NODE_COLUMN) {
172,019✔
1105
      SColumnNode*     pSColumnNode = (SColumnNode*)pNode;
172,019✔
1106
      SColumnInfoData* pColInfo = (SColumnInfoData*)taosArrayGet(pResBlock->pDataBlock, pSColumnNode->slotId);
172,019✔
1107
      QUERY_CHECK_NULL(pColInfo, code, lino, end, terrno);
172,019✔
1108
      code = colDataAssign(output.columnData, pColInfo, rows, NULL);
172,019✔
1109
    } else if (nodeType(pNode) == QUERY_NODE_VALUE) {
×
1110
      continue;
×
1111
    } else {
1112
      code = scalarCalculate(pNode, pBlockList, &output, NULL, NULL);
×
1113
    }
1114

1115
    if (code != TSDB_CODE_SUCCESS) {
172,019✔
1116
      releaseColInfoData(output.columnData);
×
1117
      goto end;
×
1118
    }
1119

1120
    void* tmp = taosArrayPush(groupData, &output.columnData);
172,019✔
1121
    QUERY_CHECK_NULL(tmp, code, lino, end, terrno);
172,019✔
1122
  }
1123

1124
  for (int i = 0; i < rows; i++) {
448,062✔
1125
    gInfo = taosArrayInit(taosArrayGetSize(groupData), sizeof(SStreamGroupValue));
341,658✔
1126
    QUERY_CHECK_NULL(gInfo, code, lino, end, terrno);
341,658✔
1127

1128
    STableKeyInfo* info = taosArrayGet(pTableListInfo->pTableList, i);
341,658✔
1129
    QUERY_CHECK_NULL(info, code, lino, end, terrno);
341,658✔
1130

1131
    for (int j = 0; j < taosArrayGetSize(groupData); j++) {
854,169✔
1132
      SColumnInfoData* pValue = (SColumnInfoData*)taosArrayGetP(groupData, j);
512,918✔
1133
        int32_t ret = buildGroupInfo(pValue, i, gInfo);
512,918✔
1134
        if (ret != TSDB_CODE_SUCCESS) {
512,511✔
1135
          qError("buildGroupInfo failed at line %d since %s", __LINE__, tstrerror(ret));
×
1136
          goto end;
×
1137
        }
1138
        if (j == tbNameIndex) {
512,511✔
1139
          SStreamGroupValue* v = taosArrayGetLast(gInfo);
341,658✔
1140
          if (v != NULL){
341,658✔
1141
            v->isTbname = true;
341,658✔
1142
            v->uid = info->uid;
341,658✔
1143
          }
1144
        }
1145
    }
1146

1147
    int32_t ret = taosHashPut(groupIdMap, &info->uid, sizeof(info->uid), &gInfo, POINTER_BYTES);
341,251✔
1148
    if (ret != TSDB_CODE_SUCCESS) {
341,658✔
1149
      qError("put groupid to map failed at line %d since %s", __LINE__, tstrerror(ret));
×
1150
      goto end;
×
1151
    }
1152
    qDebug("put groupid to map gid:%" PRIu64, info->uid);
341,658✔
1153
    gInfo = NULL;
341,658✔
1154
  }
1155

1156
end:
106,404✔
1157
  blockDataDestroy(pResBlock);
106,404✔
1158
  taosArrayDestroy(pBlockList);
106,404✔
1159
  taosArrayDestroyEx(pUidTagList, freeItem);
106,404✔
1160
  taosArrayDestroyP(groupData, releaseColInfoData);
106,404✔
1161
  taosArrayDestroyEx(gInfo, tDestroySStreamGroupValue);
106,404✔
1162

1163
  if (code != TSDB_CODE_SUCCESS) {
106,404✔
1164
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
1165
  }
1166
}
1167

1168
int32_t getColInfoResultForGroupby(void* pVnode, SNodeList* group, STableListInfo* pTableListInfo, uint8_t* digest,
1,864,471✔
1169
                                   SStorageAPI* pAPI, bool initRemainGroups, SHashObj* groupIdMap) {
1170
  int32_t      code = TSDB_CODE_SUCCESS;
1,864,471✔
1171
  int32_t      lino = 0;
1,864,471✔
1172
  SArray*      pBlockList = NULL;
1,864,471✔
1173
  SSDataBlock* pResBlock = NULL;
1,864,471✔
1174
  void*        keyBuf = NULL;
1,864,626✔
1175
  SArray*      groupData = NULL;
1,864,626✔
1176
  SArray*      pUidTagList = NULL;
1,864,626✔
1177
  SArray*      tableList = NULL;
1,864,626✔
1178
  SArray*      gInfo = NULL;
1,864,626✔
1179

1180
  int32_t rows = taosArrayGetSize(pTableListInfo->pTableList);
1,864,626✔
1181
  if (rows == 0) {
1,864,626✔
1182
    return TSDB_CODE_SUCCESS;
×
1183
  } 
1184

1185
  T_MD5_CTX context = {0};
1,864,626✔
1186
  if (tsTagFilterCache && groupIdMap == NULL) {
1,864,626✔
1187
    SNodeListNode* listNode = NULL;
×
1188
    code = nodesMakeNode(QUERY_NODE_NODE_LIST, (SNode**)&listNode);
×
1189
    if (TSDB_CODE_SUCCESS != code) {
×
1190
      goto end;
×
1191
    }
1192
    listNode->pNodeList = group;
×
1193
    code = genTbGroupDigest((SNode*)listNode, digest, &context);
×
1194
    QUERY_CHECK_CODE(code, lino, end);
×
1195

1196
    nodesFree(listNode);
×
1197

1198
    code = pAPI->metaFn.metaGetCachedTbGroup(pVnode, pTableListInfo->idInfo.suid, context.digest,
×
1199
                                             tListLen(context.digest), &tableList);
1200
    QUERY_CHECK_CODE(code, lino, end);
×
1201

1202
    if (tableList) {
×
1203
      taosArrayDestroy(pTableListInfo->pTableList);
×
1204
      pTableListInfo->pTableList = tableList;
×
1205
      qDebug("retrieve tb group list from cache, numOfTables:%d",
×
1206
             (int32_t)taosArrayGetSize(pTableListInfo->pTableList));
1207
      goto end;
×
1208
    }
1209
  }
1210

1211
  pUidTagList = taosArrayInit(8, sizeof(STUidTagInfo));
1,864,626✔
1212
  QUERY_CHECK_NULL(pUidTagList, code, lino, end, terrno);
1,864,626✔
1213

1214
  for (int32_t i = 0; i < rows; ++i) {
10,974,181✔
1215
    STableKeyInfo* pkeyInfo = taosArrayGet(pTableListInfo->pTableList, i);
9,109,555✔
1216
    QUERY_CHECK_NULL(pkeyInfo, code, lino, end, terrno);
9,109,555✔
1217
    STUidTagInfo info = {.uid = pkeyInfo->uid};
9,109,555✔
1218
    void*        tmp = taosArrayPush(pUidTagList, &info);
9,109,555✔
1219
    QUERY_CHECK_NULL(tmp, code, lino, end, terrno);
9,109,555✔
1220
  }
1221

1222
  if (taosArrayGetSize(pUidTagList) > 0) {
1,864,626✔
1223
    code = pAPI->metaFn.getTableTagsByUid(pVnode, pTableListInfo->idInfo.suid, pUidTagList);
1,864,626✔
1224
  } else {
1225
    code = pAPI->metaFn.getTableTags(pVnode, pTableListInfo->idInfo.suid, pUidTagList);
×
1226
  }
1227
  if (code != TSDB_CODE_SUCCESS) {
1,864,626✔
1228
    goto end;
×
1229
  }
1230

1231
  SArray* pColList = NULL;
1,864,626✔
1232
  code = qGetColumnsFromNodeList(group, true, &pColList); 
1,864,626✔
1233
  if (code != TSDB_CODE_SUCCESS) {
1,864,057✔
1234
    goto end;
×
1235
  }
1236

1237
  int32_t numOfTables = taosArrayGetSize(pUidTagList);
1,864,057✔
1238
  pResBlock = createTagValBlockForFilter(pColList, numOfTables, pUidTagList, pVnode, pAPI);
1,863,329✔
1239
  taosArrayDestroy(pColList);
1,863,309✔
1240
  if (pResBlock == NULL) {
1,863,534✔
1241
    code = terrno;
×
1242
    goto end;
×
1243
  }
1244

1245
  //  int64_t st1 = taosGetTimestampUs();
1246
  //  qDebug("generate tag block rows:%d, cost:%ld us", rows, st1-st);
1247

1248
  pBlockList = taosArrayInit(2, POINTER_BYTES);
1,863,534✔
1249
  QUERY_CHECK_NULL(pBlockList, code, lino, end, terrno);
1,864,626✔
1250

1251
  void* tmp = taosArrayPush(pBlockList, &pResBlock);
1,864,626✔
1252
  QUERY_CHECK_NULL(tmp, code, lino, end, terrno);
1,864,626✔
1253

1254
  groupData = taosArrayInit(2, POINTER_BYTES);
1,864,626✔
1255
  QUERY_CHECK_NULL(groupData, code, lino, end, terrno);
1,864,626✔
1256

1257
  SNode* pNode = NULL;
1,864,251✔
1258
  FOREACH(pNode, group) {
3,849,797✔
1259
    SScalarParam output = {0};
1,986,135✔
1260

1261
    switch (nodeType(pNode)) {
1,986,135✔
1262
      case QUERY_NODE_VALUE:
×
1263
        break;
×
1264
      case QUERY_NODE_COLUMN:
1,986,135✔
1265
      case QUERY_NODE_OPERATOR:
1266
      case QUERY_NODE_FUNCTION: {
1267
        SExprNode* expNode = (SExprNode*)pNode;
1,986,135✔
1268
        code = createResultData(&expNode->resType, rows, &output);
1,986,135✔
1269
        if (code != TSDB_CODE_SUCCESS) {
1,985,407✔
1270
          goto end;
×
1271
        }
1272
        break;
1,985,407✔
1273
      }
1274

1275
      default:
×
1276
        code = TSDB_CODE_OPS_NOT_SUPPORT;
×
1277
        goto end;
×
1278
    }
1279

1280
    if (nodeType(pNode) == QUERY_NODE_COLUMN) {
1,985,407✔
1281
      SColumnNode*     pSColumnNode = (SColumnNode*)pNode;
1,975,970✔
1282
      SColumnInfoData* pColInfo = (SColumnInfoData*)taosArrayGet(pResBlock->pDataBlock, pSColumnNode->slotId);
1,975,970✔
1283
      QUERY_CHECK_NULL(pColInfo, code, lino, end, terrno);
1,975,757✔
1284
      code = colDataAssign(output.columnData, pColInfo, rows, NULL);
1,975,757✔
1285
    } else if (nodeType(pNode) == QUERY_NODE_VALUE) {
10,165✔
1286
      continue;
×
1287
    } else {
1288
      code = scalarCalculate(pNode, pBlockList, &output, NULL, NULL);
10,165✔
1289
    }
1290

1291
    if (code != TSDB_CODE_SUCCESS) {
1,984,090✔
1292
      releaseColInfoData(output.columnData);
×
1293
      goto end;
×
1294
    }
1295

1296
    void* tmp = taosArrayPush(groupData, &output.columnData);
1,985,771✔
1297
    QUERY_CHECK_NULL(tmp, code, lino, end, terrno);
1,985,771✔
1298
  }
1299

1300
  int32_t keyLen = 0;
1,864,037✔
1301
  SNode*  node;
1302
  FOREACH(node, group) {
3,848,857✔
1303
    SExprNode* pExpr = (SExprNode*)node;
1,984,594✔
1304
    keyLen += pExpr->resType.bytes;
1,984,594✔
1305
  }
1306

1307
  int32_t nullFlagSize = sizeof(int8_t) * LIST_LENGTH(group);
1,863,310✔
1308
  keyLen += nullFlagSize;
1,862,935✔
1309

1310
  keyBuf = taosMemoryCalloc(1, keyLen);
1,862,935✔
1311
  if (keyBuf == NULL) {
1,864,626✔
1312
    code = terrno;
×
1313
    goto end;
×
1314
  }
1315

1316
  if (initRemainGroups) {
1,864,626✔
1317
    pTableListInfo->remainGroups =
850,235✔
1318
        taosHashInit(rows, taosGetDefaultHashFunction(TSDB_DATA_TYPE_BIGINT), false, HASH_NO_LOCK);
850,235✔
1319
    if (pTableListInfo->remainGroups == NULL) {
850,235✔
1320
      code = terrno;
×
1321
      goto end;
×
1322
    }
1323
  }
1324

1325
  for (int i = 0; i < rows; i++) {
10,972,046✔
1326
    STableKeyInfo* info = taosArrayGet(pTableListInfo->pTableList, i);
9,107,570✔
1327
    QUERY_CHECK_NULL(info, code, lino, end, terrno);
9,107,566✔
1328

1329
    if (groupIdMap != NULL){
9,107,566✔
1330
      gInfo = taosArrayInit(taosArrayGetSize(groupData), sizeof(SStreamGroupValue));
176,674✔
1331
    }
1332
    
1333
    char* isNull = (char*)keyBuf;
9,108,523✔
1334
    char* pStart = (char*)keyBuf + sizeof(int8_t) * LIST_LENGTH(group);
9,108,523✔
1335
    for (int j = 0; j < taosArrayGetSize(groupData); j++) {
18,860,725✔
1336
      SColumnInfoData* pValue = (SColumnInfoData*)taosArrayGetP(groupData, j);
9,751,746✔
1337

1338
      if (groupIdMap != NULL && gInfo != NULL) {
9,751,908✔
1339
        int32_t ret = buildGroupInfo(pValue, i, gInfo);
201,038✔
1340
        if (ret != TSDB_CODE_SUCCESS) {
201,402✔
1341
          qError("buildGroupInfo failed at line %d since %s", __LINE__, tstrerror(ret));
×
1342
          taosArrayDestroyEx(gInfo, tDestroySStreamGroupValue);
×
1343
          gInfo = NULL;
×
1344
        }
1345
      }
1346
      
1347
      if (colDataIsNull_s(pValue, i)) {
19,505,133✔
1348
        isNull[j] = 1;
94,705✔
1349
      } else {
1350
        isNull[j] = 0;
9,658,156✔
1351
        char* data = colDataGetData(pValue, i);
9,657,475✔
1352
        if (pValue->info.type == TSDB_DATA_TYPE_JSON) {
9,658,413✔
1353
          // if (tTagIsJson(data)) {
1354
          //   code = TSDB_CODE_QRY_JSON_IN_GROUP_ERROR;
1355
          //   goto end;
1356
          // }
1357
          if (tTagIsJsonNull(data)) {
89,487✔
1358
            isNull[j] = 1;
×
1359
            continue;
×
1360
          }
1361
          int32_t len = getJsonValueLen(data);
89,487✔
1362
          memcpy(pStart, data, len);
89,487✔
1363
          pStart += len;
89,487✔
1364
        } else if (IS_VAR_DATA_TYPE(pValue->info.type)) {
9,568,177✔
1365
          if (IS_STR_DATA_BLOB(pValue->info.type)) {
6,572,171✔
1366
            if (blobDataTLen(data) > TSDB_MAX_BLOB_LEN) {
465✔
1367
              code = TSDB_CODE_TDB_INVALID_TABLE_SCHEMA_VER;
×
1368
              goto end;
×
1369
            }
1370
            memcpy(pStart, data, blobDataTLen(data));
×
1371
            pStart += blobDataTLen(data);
×
1372
          } else {
1373
            if (varDataTLen(data) > pValue->info.bytes) {
6,572,639✔
1374
              code = TSDB_CODE_TDB_INVALID_TABLE_SCHEMA_VER;
×
1375
              goto end;
×
1376
            }
1377
            memcpy(pStart, data, varDataTLen(data));
6,572,634✔
1378
            pStart += varDataTLen(data);
6,572,634✔
1379
          }
1380
        } else {
1381
          memcpy(pStart, data, pValue->info.bytes);
2,996,340✔
1382
          pStart += pValue->info.bytes;
2,995,751✔
1383
        }
1384
      }
1385
    }
1386

1387
    int32_t len = (int32_t)(pStart - (char*)keyBuf);
9,107,279✔
1388
    info->groupId = calcGroupId(keyBuf, len);
9,107,279✔
1389
    if (groupIdMap != NULL && gInfo != NULL) {
9,108,940✔
1390
      int32_t ret = taosHashPut(groupIdMap, &info->groupId, sizeof(info->groupId), &gInfo, POINTER_BYTES);
177,777✔
1391
      if (ret != TSDB_CODE_SUCCESS) {
177,777✔
1392
        qError("put groupid to map failed at line %d since %s", __LINE__, tstrerror(ret));
×
1393
        taosArrayDestroyEx(gInfo, tDestroySStreamGroupValue);
×
1394
      }
1395
      qDebug("put groupid to map gid:%" PRIu64, info->groupId);
177,777✔
1396
      gInfo = NULL;
177,777✔
1397
    }
1398
    if (initRemainGroups) {
9,108,940✔
1399
      // groupId ~ table uid
1400
      code = taosHashPut(pTableListInfo->remainGroups, &(info->groupId), sizeof(info->groupId), &(info->uid),
4,555,619✔
1401
                         sizeof(info->uid));
1402
      if (code == TSDB_CODE_DUP_KEY) {
4,555,566✔
1403
        code = TSDB_CODE_SUCCESS;
829,747✔
1404
      }
1405
      QUERY_CHECK_CODE(code, lino, end);
4,555,566✔
1406
    }
1407
  }
1408

1409
  if (tsTagFilterCache && groupIdMap == NULL) {
1,864,476✔
1410
    tableList = taosArrayDup(pTableListInfo->pTableList, NULL);
×
1411
    QUERY_CHECK_NULL(tableList, code, lino, end, terrno);
×
1412

1413
    code = pAPI->metaFn.metaPutTbGroupToCache(pVnode, pTableListInfo->idInfo.suid, context.digest,
×
1414
                                              tListLen(context.digest), tableList,
1415
                                              taosArrayGetSize(tableList) * sizeof(STableKeyInfo));
×
1416
    QUERY_CHECK_CODE(code, lino, end);
×
1417
  }
1418

1419
  //  int64_t st2 = taosGetTimestampUs();
1420
  //  qDebug("calculate tag block rows:%d, cost:%ld us", rows, st2-st1);
1421

1422
end:
1,863,748✔
1423
  taosMemoryFreeClear(keyBuf);
1,864,626✔
1424
  blockDataDestroy(pResBlock);
1,864,626✔
1425
  taosArrayDestroy(pBlockList);
1,864,626✔
1426
  taosArrayDestroyEx(pUidTagList, freeItem);
1,864,626✔
1427
  taosArrayDestroyP(groupData, releaseColInfoData);
1,864,626✔
1428
  taosArrayDestroyEx(gInfo, tDestroySStreamGroupValue);
1,864,626✔
1429

1430
  if (code != TSDB_CODE_SUCCESS) {
1,864,626✔
1431
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
1432
  }
1433
  return code;
1,864,626✔
1434
}
1435

1436
static int32_t nameComparFn(const void* p1, const void* p2) {
710,398✔
1437
  const char* pName1 = *(const char**)p1;
710,398✔
1438
  const char* pName2 = *(const char**)p2;
710,398✔
1439

1440
  int32_t ret = strcmp(pName1, pName2);
710,398✔
1441
  if (ret == 0) {
710,398✔
1442
    return 0;
18,210✔
1443
  } else {
1444
    return (ret > 0) ? 1 : -1;
692,188✔
1445
  }
1446
}
1447

1448
static SArray* getTableNameList(const SNodeListNode* pList) {
419,688✔
1449
  int32_t    code = TSDB_CODE_SUCCESS;
419,688✔
1450
  int32_t    lino = 0;
419,688✔
1451
  int32_t    len = LIST_LENGTH(pList->pNodeList);
419,688✔
1452
  SListCell* cell = pList->pNodeList->pHead;
419,688✔
1453

1454
  SArray* pTbList = taosArrayInit(len, POINTER_BYTES);
419,688✔
1455
  QUERY_CHECK_NULL(pTbList, code, lino, _end, terrno);
419,688✔
1456

1457
  for (int i = 0; i < pList->pNodeList->length; i++) {
1,140,339✔
1458
    SValueNode* valueNode = (SValueNode*)cell->pNode;
720,651✔
1459
    if (!IS_VAR_DATA_TYPE(valueNode->node.resType.type)) {
720,651✔
1460
      terrno = TSDB_CODE_INVALID_PARA;
×
1461
      taosArrayDestroy(pTbList);
×
1462
      return NULL;
×
1463
    }
1464

1465
    char* name = varDataVal(valueNode->datum.p);
720,651✔
1466
    void* tmp = taosArrayPush(pTbList, &name);
720,651✔
1467
    QUERY_CHECK_NULL(tmp, code, lino, _end, terrno);
720,651✔
1468
    cell = cell->pNext;
720,651✔
1469
  }
1470

1471
  size_t numOfTables = taosArrayGetSize(pTbList);
419,688✔
1472

1473
  // order the name
1474
  taosArraySort(pTbList, nameComparFn);
419,688✔
1475

1476
  // remove the duplicates
1477
  SArray* pNewList = taosArrayInit(taosArrayGetSize(pTbList), sizeof(void*));
419,688✔
1478
  QUERY_CHECK_NULL(pNewList, code, lino, _end, terrno);
419,688✔
1479
  void* tmpTbl = taosArrayGet(pTbList, 0);
419,688✔
1480
  QUERY_CHECK_NULL(tmpTbl, code, lino, _end, terrno);
419,688✔
1481
  void* tmp = taosArrayPush(pNewList, tmpTbl);
419,688✔
1482
  QUERY_CHECK_NULL(tmp, code, lino, _end, terrno);
419,688✔
1483

1484
  for (int32_t i = 1; i < numOfTables; ++i) {
720,651✔
1485
    char** name = taosArrayGetLast(pNewList);
300,963✔
1486
    char** nameInOldList = taosArrayGet(pTbList, i);
300,963✔
1487
    QUERY_CHECK_NULL(nameInOldList, code, lino, _end, terrno);
300,963✔
1488
    if (strcmp(*name, *nameInOldList) == 0) {
300,963✔
1489
      continue;
9,786✔
1490
    }
1491

1492
    tmp = taosArrayPush(pNewList, nameInOldList);
291,177✔
1493
    QUERY_CHECK_NULL(tmp, code, lino, _end, terrno);
291,177✔
1494
  }
1495

1496
_end:
419,688✔
1497
  taosArrayDestroy(pTbList);
419,688✔
1498
  if (code != TSDB_CODE_SUCCESS) {
419,688✔
1499
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
1500
    return NULL;
×
1501
  }
1502
  return pNewList;
419,688✔
1503
}
1504

1505
static int tableUidCompare(const void* a, const void* b) {
×
1506
  uint64_t u1 = *(uint64_t*)a;
×
1507
  uint64_t u2 = *(uint64_t*)b;
×
1508

1509
  if (u1 == u2) {
×
1510
    return 0;
×
1511
  }
1512

1513
  return u1 < u2 ? -1 : 1;
×
1514
}
1515

1516
static int32_t filterTableInfoCompare(const void* a, const void* b) {
18,392,111✔
1517
  STUidTagInfo* p1 = (STUidTagInfo*)a;
18,392,111✔
1518
  STUidTagInfo* p2 = (STUidTagInfo*)b;
18,392,111✔
1519

1520
  if (p1->uid == p2->uid) {
18,392,111✔
1521
    return 0;
×
1522
  }
1523

1524
  return p1->uid < p2->uid ? -1 : 1;
18,392,111✔
1525
}
1526

1527
static FilterCondType checkTagCond(SNode* cond) {
12,878,887✔
1528
  if (nodeType(cond) == QUERY_NODE_OPERATOR) {
12,878,887✔
1529
    return FILTER_NO_LOGIC;
10,826,738✔
1530
  }
1531
  if (nodeType(cond) != QUERY_NODE_LOGIC_CONDITION || ((SLogicConditionNode*)cond)->condType != LOGIC_COND_TYPE_AND) {
2,052,528✔
1532
    return FILTER_AND;
215,661✔
1533
  }
1534
  return FILTER_OTHER;
1,836,870✔
1535
}
1536

1537
static int32_t optimizeTbnameInCond(void* pVnode, int64_t suid, SArray* list, SNode* cond, SStorageAPI* pAPI) {
12,879,266✔
1538
  int32_t ret = -1;
12,879,266✔
1539
  int32_t ntype = nodeType(cond);
12,879,266✔
1540

1541
  if (ntype == QUERY_NODE_OPERATOR) {
12,879,663✔
1542
    ret = optimizeTbnameInCondImpl(pVnode, list, cond, pAPI, suid);
10,826,738✔
1543
  }
1544

1545
  if (ntype != QUERY_NODE_LOGIC_CONDITION || ((SLogicConditionNode*)cond)->condType != LOGIC_COND_TYPE_AND) {
12,879,266✔
1546
    return ret;
11,042,199✔
1547
  }
1548

1549
  bool                 hasTbnameCond = false;
1,837,067✔
1550
  SLogicConditionNode* pNode = (SLogicConditionNode*)cond;
1,837,067✔
1551
  SNodeList*           pList = (SNodeList*)pNode->pParameterList;
1,837,067✔
1552

1553
  int32_t len = LIST_LENGTH(pList);
1,837,067✔
1554
  if (len <= 0) {
1,836,867✔
1555
    return ret;
×
1556
  }
1557

1558
  SListCell* cell = pList->pHead;
1,836,867✔
1559
  for (int i = 0; i < len; i++) {
5,925,907✔
1560
    if (cell == NULL) break;
4,095,260✔
1561
    if (optimizeTbnameInCondImpl(pVnode, list, cell->pNode, pAPI, suid) == 0) {
4,095,260✔
1562
      hasTbnameCond = true;
6,223✔
1563
      break;
6,223✔
1564
    }
1565
    cell = cell->pNext;
4,088,840✔
1566
  }
1567

1568
  taosArraySort(list, filterTableInfoCompare);
1,836,870✔
1569
  taosArrayRemoveDuplicate(list, filterTableInfoCompare, NULL);
1,836,867✔
1570

1571
  if (hasTbnameCond) {
1,836,467✔
1572
    ret = pAPI->metaFn.getTableTagsByUid(pVnode, suid, list);
6,223✔
1573
  }
1574

1575
  return ret;
1,836,867✔
1576
}
1577

1578
// only return uid that does not contained in pExistedUidList
1579
static int32_t optimizeTbnameInCondImpl(void* pVnode, SArray* pExistedUidList, SNode* pTagCond, SStorageAPI* pStoreAPI,
14,922,198✔
1580
                                        uint64_t suid) {
1581
  if (nodeType(pTagCond) != QUERY_NODE_OPERATOR) {
14,922,198✔
1582
    return -1;
6,248✔
1583
  }
1584

1585
  SOperatorNode* pNode = (SOperatorNode*)pTagCond;
14,915,950✔
1586
  if (pNode->opType != OP_TYPE_IN) {
14,915,950✔
1587
    return -1;
14,185,615✔
1588
  }
1589

1590
  if ((pNode->pLeft != NULL && ((nodeType(pNode->pLeft) == QUERY_NODE_FUNCTION &&
730,135✔
1591
                                 ((SFunctionNode*)pNode->pLeft)->funcType == FUNCTION_TYPE_TBNAME)) ||
419,688✔
1592
       (nodeType(pNode->pLeft) == QUERY_NODE_COLUMN && ((SColumnNode*)pNode->pLeft)->colType == COLUMN_TYPE_TBNAME)) &&
310,447✔
1593
      (pNode->pRight != NULL && nodeType(pNode->pRight) == QUERY_NODE_NODE_LIST)) {
419,688✔
1594
    SNodeListNode* pList = (SNodeListNode*)pNode->pRight;
419,688✔
1595

1596
    int32_t len = LIST_LENGTH(pList->pNodeList);
419,688✔
1597
    if (len <= 0) {
419,688✔
1598
      return -1;
×
1599
    }
1600

1601
    SArray*   pTbList = getTableNameList(pList);
419,688✔
1602
    int32_t   numOfTables = taosArrayGetSize(pTbList);
419,688✔
1603
    SHashObj* uHash = NULL;
419,688✔
1604

1605
    size_t numOfExisted = taosArrayGetSize(pExistedUidList);  // len > 0 means there already have uids
419,688✔
1606
    if (numOfExisted > 0) {
419,688✔
1607
      uHash = taosHashInit(numOfExisted / 0.7, taosGetDefaultHashFunction(TSDB_DATA_TYPE_BIGINT), false, HASH_NO_LOCK);
2,354✔
1608
      if (!uHash) {
2,354✔
1609
        qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
1610
        return terrno;
×
1611
      }
1612

1613
      for (int i = 0; i < numOfExisted; i++) {
2,352,823✔
1614
        STUidTagInfo* pTInfo = taosArrayGet(pExistedUidList, i);
2,350,469✔
1615
        if (!pTInfo) {
2,350,469✔
1616
          qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
1617
          return terrno;
×
1618
        }
1619
        int32_t tempRes = taosHashPut(uHash, &pTInfo->uid, sizeof(uint64_t), &i, sizeof(i));
2,350,469✔
1620
        if (tempRes != TSDB_CODE_SUCCESS && tempRes != TSDB_CODE_DUP_KEY) {
2,350,469✔
1621
          qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(tempRes));
×
1622
          return tempRes;
×
1623
        }
1624
      }
1625
    }
1626

1627
    for (int i = 0; i < numOfTables; i++) {
1,113,682✔
1628
      char* name = taosArrayGetP(pTbList, i);
700,490✔
1629

1630
      uint64_t uid = 0, csuid = 0;
700,490✔
1631
      if (pStoreAPI->metaFn.getTableUidByName(pVnode, name, &uid) == 0) {
700,490✔
1632
        ETableType tbType = TSDB_TABLE_MAX;
419,814✔
1633
        if (pStoreAPI->metaFn.getTableTypeSuidByName(pVnode, name, &tbType, &csuid) == 0 &&
419,814✔
1634
            tbType == TSDB_CHILD_TABLE) {
419,814✔
1635
          if (suid != csuid) {
413,318✔
1636
            continue;
852✔
1637
          }
1638
          if (NULL == uHash || taosHashGet(uHash, &uid, sizeof(uid)) == NULL) {
412,466✔
1639
            STUidTagInfo s = {.uid = uid, .name = name, .pTagVal = NULL};
411,289✔
1640
            void*        tmp = taosArrayPush(pExistedUidList, &s);
411,289✔
1641
            if (!tmp) {
411,289✔
1642
              return terrno;
×
1643
            }
1644
          }
1645
        } else {
1646
          taosArrayDestroy(pTbList);
6,496✔
1647
          taosHashCleanup(uHash);
6,496✔
1648
          return -1;
6,496✔
1649
        }
1650
      } else {
1651
        //        qWarn("failed to get tableIds from by table name: %s, reason: %s", name, tstrerror(terrno));
1652
        terrno = 0;
280,676✔
1653
      }
1654
    }
1655

1656
    taosHashCleanup(uHash);
413,192✔
1657
    taosArrayDestroy(pTbList);
413,192✔
1658
    return 0;
413,192✔
1659
  }
1660

1661
  return -1;
310,447✔
1662
}
1663

1664
SSDataBlock* createTagValBlockForFilter(SArray* pColList, int32_t numOfTables, SArray* pUidTagList, void* pVnode,
13,729,461✔
1665
                                        SStorageAPI* pStorageAPI) {
1666
  int32_t      code = TSDB_CODE_SUCCESS;
13,729,461✔
1667
  int32_t      lino = 0;
13,729,461✔
1668
  SSDataBlock* pResBlock = NULL;
13,729,461✔
1669
  code = createDataBlock(&pResBlock);
13,730,295✔
1670
  QUERY_CHECK_CODE(code, lino, _end);
13,730,266✔
1671

1672
  for (int32_t i = 0; i < taosArrayGetSize(pColList); ++i) {
28,785,449✔
1673
    SColumnInfoData colInfo = {0};
15,055,311✔
1674
    void*           tmp = taosArrayGet(pColList, i);
15,055,191✔
1675
    QUERY_CHECK_NULL(tmp, code, lino, _end, terrno);
15,055,312✔
1676
    colInfo.info = *(SColumnInfo*)tmp;
15,055,312✔
1677
    code = blockDataAppendColInfo(pResBlock, &colInfo);
15,054,166✔
1678
    QUERY_CHECK_CODE(code, lino, _end);
15,055,708✔
1679
  }
1680

1681
  code = blockDataEnsureCapacity(pResBlock, numOfTables);
13,730,016✔
1682
  if (code != TSDB_CODE_SUCCESS) {
13,728,497✔
1683
    terrno = code;
×
1684
    blockDataDestroy(pResBlock);
×
1685
    return NULL;
×
1686
  }
1687

1688
  pResBlock->info.rows = numOfTables;
13,728,497✔
1689

1690
  int32_t numOfCols = taosArrayGetSize(pResBlock->pDataBlock);
13,728,846✔
1691

1692
  for (int32_t i = 0; i < numOfTables; i++) {
201,103,611✔
1693
    STUidTagInfo* p1 = taosArrayGet(pUidTagList, i);
187,372,740✔
1694
    QUERY_CHECK_NULL(p1, code, lino, _end, terrno);
187,379,366✔
1695

1696
    for (int32_t j = 0; j < numOfCols; j++) {
382,561,317✔
1697
      SColumnInfoData* pColInfo = (SColumnInfoData*)taosArrayGet(pResBlock->pDataBlock, j);
195,176,108✔
1698
      QUERY_CHECK_NULL(pColInfo, code, lino, _end, terrno);
195,162,254✔
1699

1700
      if (pColInfo->info.colId == -1) {  // tbname
195,162,254✔
1701
        char str[TSDB_TABLE_FNAME_LEN + VARSTR_HEADER_SIZE] = {0};
8,011,626✔
1702
        if (p1->name != NULL) {
8,012,806✔
1703
          STR_TO_VARSTR(str, p1->name);
411,289✔
1704
        } else {  // name is not retrieved during filter
1705
          code = pStorageAPI->metaFn.getTableNameByUid(pVnode, p1->uid, str);
7,600,525✔
1706
          QUERY_CHECK_CODE(code, lino, _end);
7,593,291✔
1707
        }
1708

1709
        code = colDataSetVal(pColInfo, i, str, false);
8,004,580✔
1710
        QUERY_CHECK_CODE(code, lino, _end);
7,984,585✔
1711
#if TAG_FILTER_DEBUG
1712
        qDebug("tagfilter uid:%ld, tbname:%s", *uid, str + 2);
1713
#endif
1714
      } else {
1715
        STagVal tagVal = {0};
187,159,703✔
1716
        tagVal.cid = pColInfo->info.colId;
187,164,994✔
1717
        if (p1->pTagVal == NULL) {
187,171,297✔
1718
          colDataSetNULL(pColInfo, i);
9,100✔
1719
        } else {
1720
          const char* p = pStorageAPI->metaFn.extractTagVal(p1->pTagVal, pColInfo->info.type, &tagVal);
187,153,199✔
1721

1722
          if (p == NULL || (pColInfo->info.type == TSDB_DATA_TYPE_JSON && ((STag*)p)->nTag == 0)) {
187,171,648✔
1723
            colDataSetNULL(pColInfo, i);
3,812,312✔
1724
          } else if (pColInfo->info.type == TSDB_DATA_TYPE_JSON) {
183,353,332✔
1725
            code = colDataSetVal(pColInfo, i, p, false);
707,812✔
1726
            QUERY_CHECK_CODE(code, lino, _end);
707,812✔
1727
          } else if (IS_VAR_DATA_TYPE(pColInfo->info.type)) {
295,410,367✔
1728
            if (IS_STR_DATA_BLOB(pColInfo->info.type)) {
112,752,984✔
1729
              QUERY_CHECK_CODE(code = TSDB_CODE_BLOB_NOT_SUPPORT_TAG, lino, _end);
×
1730
            }
1731
            char* tmp = taosMemoryMalloc(tagVal.nData + VARSTR_HEADER_SIZE + 1);
112,775,149✔
1732
            QUERY_CHECK_NULL(tmp, code, lino, _end, terrno);
112,780,199✔
1733
            varDataSetLen(tmp, tagVal.nData);
112,780,199✔
1734
            memcpy(tmp + VARSTR_HEADER_SIZE, tagVal.pData, tagVal.nData);
112,781,013✔
1735
            code = colDataSetVal(pColInfo, i, tmp, false);
112,779,888✔
1736
#if TAG_FILTER_DEBUG
1737
            qDebug("tagfilter varch:%s", tmp + 2);
1738
#endif
1739
            taosMemoryFree(tmp);
112,781,724✔
1740
            QUERY_CHECK_CODE(code, lino, _end);
112,782,670✔
1741
          } else {
1742
            code = colDataSetVal(pColInfo, i, (const char*)&tagVal.i64, false);
69,863,264✔
1743
            QUERY_CHECK_CODE(code, lino, _end);
69,884,222✔
1744
#if TAG_FILTER_DEBUG
1745
            if (pColInfo->info.type == TSDB_DATA_TYPE_INT) {
1746
              qDebug("tagfilter int:%d", *(int*)(&tagVal.i64));
1747
            } else if (pColInfo->info.type == TSDB_DATA_TYPE_DOUBLE) {
1748
              qDebug("tagfilter double:%f", *(double*)(&tagVal.i64));
1749
            }
1750
#endif
1751
          }
1752
        }
1753
      }
1754
    }
1755
  }
1756

1757
_end:
13,717,000✔
1758
  if (code != TSDB_CODE_SUCCESS) {
13,730,871✔
1759
    blockDataDestroy(pResBlock);
875✔
1760
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
1761
    terrno = code;
×
1762
    return NULL;
×
1763
  }
1764
  return pResBlock;
13,729,996✔
1765
}
1766

1767
static int32_t doSetQualifiedUid(STableListInfo* pListInfo, SArray* pUidList, const SArray* pUidTagList,
11,086,630✔
1768
                                 bool* pResultList, bool addUid) {
1769
  taosArrayClear(pUidList);
11,086,630✔
1770

1771
  STableKeyInfo info = {.uid = 0, .groupId = 0};
11,087,892✔
1772
  int32_t       numOfTables = taosArrayGetSize(pUidTagList);
11,088,766✔
1773
  for (int32_t i = 0; i < numOfTables; ++i) {
188,031,398✔
1774
    if (pResultList[i]) {
176,939,954✔
1775
      STUidTagInfo* tmpTag = (STUidTagInfo*)taosArrayGet(pUidTagList, i);
77,948,164✔
1776
      if (!tmpTag) {
77,946,208✔
1777
        qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
1778
        return terrno;
×
1779
      }
1780
      uint64_t uid = tmpTag->uid;
77,946,208✔
1781
      qDebug("tagfilter get uid:%" PRId64 ", res:%d", uid, pResultList[i]);
77,943,879✔
1782

1783
      info.uid = uid;
77,951,590✔
1784
      //qInfo("doSetQualifiedUid row:%d added to pTableList", i);
1785
      void* p = taosArrayPush(pListInfo->pTableList, &info);
77,951,590✔
1786
      if (p == NULL) {
77,952,414✔
1787
        return terrno;
×
1788
      }
1789

1790
      if (addUid) {
77,952,414✔
1791
        //qInfo("doSetQualifiedUid row:%d added to pUidList", i);
1792
        void* tmp = taosArrayPush(pUidList, &uid);
20,449✔
1793
        if (tmp == NULL) {
20,449✔
1794
          return terrno;
×
1795
        }
1796
      }
1797
    } else {
1798
      //qInfo("doSetQualifiedUid row:%d failed", i);
1799
    }
1800
  }
1801

1802
  return TSDB_CODE_SUCCESS;
11,091,444✔
1803
}
1804

1805
static int32_t copyExistedUids(SArray* pUidTagList, const SArray* pUidList) {
12,879,266✔
1806
  int32_t code = TSDB_CODE_SUCCESS;
12,879,266✔
1807
  int32_t numOfExisted = taosArrayGetSize(pUidList);
12,879,266✔
1808
  if (numOfExisted == 0) {
12,879,066✔
1809
    return code;
9,986,829✔
1810
  }
1811

1812
  for (int32_t i = 0; i < numOfExisted; ++i) {
35,386,228✔
1813
    uint64_t* uid = taosArrayGet(pUidList, i);
32,494,905✔
1814
    if (!uid) {
32,494,448✔
1815
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
1816
      return terrno;
×
1817
    }
1818
    STUidTagInfo info = {.uid = *uid};
32,494,448✔
1819
    void*        tmp = taosArrayPush(pUidTagList, &info);
32,493,991✔
1820
    if (!tmp) {
32,493,991✔
1821
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
1822
      return code;
×
1823
    }
1824
  }
1825
  return code;
2,891,323✔
1826
}
1827

1828
int32_t doFilterByTagCond(STableListInfo* pListInfo, SArray* pUidList, SNode* pTagCond, void* pVnode,
118,729,508✔
1829
                                 SIdxFltStatus status, SStorageAPI* pAPI, bool addUid, bool* listAdded, void* pStreamInfo) {
1830
  *listAdded = false;
118,729,508✔
1831
  if (pTagCond == NULL) {
118,748,156✔
1832
    return TSDB_CODE_SUCCESS;
105,847,476✔
1833
  }
1834

1835
  terrno = TSDB_CODE_SUCCESS;
12,900,680✔
1836

1837
  int32_t      lino = 0;
12,878,669✔
1838
  int32_t      code = TSDB_CODE_SUCCESS;
12,878,669✔
1839
  SArray*      pBlockList = NULL;
12,878,669✔
1840
  SSDataBlock* pResBlock = NULL;
12,878,669✔
1841
  SScalarParam output = {0};
12,878,272✔
1842
  SArray*      pUidTagList = NULL;
12,878,272✔
1843

1844
  SDataType type = {.type = TSDB_DATA_TYPE_BOOL, .bytes = sizeof(bool)};
12,878,272✔
1845

1846
  //  int64_t stt = taosGetTimestampUs();
1847
  pUidTagList = taosArrayInit(10, sizeof(STUidTagInfo));
12,879,069✔
1848
  QUERY_CHECK_NULL(pUidTagList, code, lino, end, terrno);
12,879,066✔
1849

1850
  code = copyExistedUids(pUidTagList, pUidList);
12,879,066✔
1851
  QUERY_CHECK_CODE(code, lino, end);
12,878,869✔
1852

1853
  FilterCondType condType = checkTagCond(pTagCond);
12,878,869✔
1854

1855
  int32_t filter = optimizeTbnameInCond(pVnode, pListInfo->idInfo.suid, pUidTagList, pTagCond, pAPI);
12,879,072✔
1856
  if (filter == 0) {  // tbname in filter is activated, do nothing and return
12,878,669✔
1857
    taosArrayClear(pUidList);
413,192✔
1858

1859
    int32_t numOfRows = taosArrayGetSize(pUidTagList);
413,192✔
1860
    code = taosArrayEnsureCap(pUidList, numOfRows);
413,192✔
1861
    QUERY_CHECK_CODE(code, lino, end);
413,192✔
1862

1863
    for (int32_t i = 0; i < numOfRows; ++i) {
3,174,950✔
1864
      STUidTagInfo* pInfo = taosArrayGet(pUidTagList, i);
2,761,758✔
1865
      QUERY_CHECK_NULL(pInfo, code, lino, end, terrno);
2,761,758✔
1866
      void* tmp = taosArrayPush(pUidList, &pInfo->uid);
2,761,758✔
1867
      QUERY_CHECK_NULL(tmp, code, lino, end, terrno);
2,761,758✔
1868
    }
1869
    terrno = 0;
413,192✔
1870
  } else {
1871
    qDebug("pUidTagList size:%d", (int32_t)taosArrayGetSize(pUidTagList));
12,465,477✔
1872
    
1873
    if (((condType == FILTER_NO_LOGIC || condType == FILTER_AND) && status != SFLT_NOT_INDEX) ||
21,865,176✔
1874
          taosArrayGetSize(pUidTagList) > 0) {
9,399,502✔
1875
      code = pAPI->metaFn.getTableTagsByUid(pVnode, pListInfo->idInfo.suid, pUidTagList);
3,560,005✔
1876
    } else {
1877
      code = pAPI->metaFn.getTableTags(pVnode, pListInfo->idInfo.suid, pUidTagList);
8,905,669✔
1878
    }
1879
    if (code != TSDB_CODE_SUCCESS) {
12,465,508✔
1880
      qError("failed to get table tags from meta, reason:%s, suid:%" PRIu64, tstrerror(code), pListInfo->idInfo.suid);
×
1881
      terrno = code;
×
1882
      QUERY_CHECK_CODE(code, lino, end);
×
1883
    }
1884
  }
1885

1886
  qDebug("final pUidTagList size:%d", (int32_t)taosArrayGetSize(pUidTagList));
12,878,700✔
1887

1888
  int32_t numOfTables = taosArrayGetSize(pUidTagList);
12,879,220✔
1889
  if (numOfTables == 0) {
12,879,546✔
1890
    goto end;
1,785,818✔
1891
  }
1892

1893
  SArray* pColList = NULL;
11,093,728✔
1894
  code = qGetColumnsFromNodeList(pTagCond, false, &pColList); 
11,093,728✔
1895
  if (code != TSDB_CODE_SUCCESS) {
11,092,820✔
1896
    goto end;
×
1897
  }
1898
  pResBlock = createTagValBlockForFilter(pColList, numOfTables, pUidTagList, pVnode, pAPI);
11,092,820✔
1899
  taosArrayDestroy(pColList);
11,092,891✔
1900
  if (pResBlock == NULL) {
11,092,805✔
1901
    code = terrno;
×
1902
    QUERY_CHECK_CODE(code, lino, end);
×
1903
  }
1904

1905
  //fprintDataBlock(pResBlock, "tagFilter", "", 0);
1906

1907
  //  int64_t st1 = taosGetTimestampUs();
1908
  //  qDebug("generate tag block rows:%d, cost:%ld us", rows, st1-st);
1909
  pBlockList = taosArrayInit(2, POINTER_BYTES);
11,092,805✔
1910
  QUERY_CHECK_NULL(pBlockList, code, lino, end, terrno);
11,091,057✔
1911

1912
  void* tmp = taosArrayPush(pBlockList, &pResBlock);
11,092,735✔
1913
  QUERY_CHECK_NULL(tmp, code, lino, end, terrno);
11,092,735✔
1914

1915
  code = createResultData(&type, numOfTables, &output);
11,092,735✔
1916
  if (code != TSDB_CODE_SUCCESS) {
11,091,926✔
1917
    terrno = code;
×
1918
    QUERY_CHECK_CODE(code, lino, end);
×
1919
  }
1920

1921
  code = scalarCalculate(pTagCond, pBlockList, &output, pStreamInfo, NULL);
11,091,926✔
1922
  if (code != TSDB_CODE_SUCCESS) {
11,087,363✔
1923
    qError("failed to calculate scalar, reason:%s", tstrerror(code));
1,098✔
1924
    terrno = code;
1,098✔
1925
    QUERY_CHECK_CODE(code, lino, end);
1,098✔
1926
  }
1927

1928
  code = doSetQualifiedUid(pListInfo, pUidList, pUidTagList, (bool*)output.columnData->pData, addUid);
11,086,265✔
1929
  if (code != TSDB_CODE_SUCCESS) {
11,091,641✔
1930
    terrno = code;
×
1931
    QUERY_CHECK_CODE(code, lino, end);
×
1932
  }
1933
  *listAdded = true;
11,091,641✔
1934

1935
end:
12,878,528✔
1936
  if (code != TSDB_CODE_SUCCESS) {
12,877,984✔
1937
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
1,098✔
1938
  }
1939
  blockDataDestroy(pResBlock);
12,877,984✔
1940
  taosArrayDestroy(pBlockList);
12,878,545✔
1941
  taosArrayDestroyEx(pUidTagList, freeItem);
12,878,850✔
1942

1943
  colDataDestroy(output.columnData);
12,879,057✔
1944
  taosMemoryFreeClear(output.columnData);
12,879,195✔
1945
  return code;
12,879,069✔
1946
}
1947

1948
typedef struct {
1949
  int32_t code;
1950
  SStreamRuntimeFuncInfo* pStreamRuntimeInfo;
1951
} PlaceHolderContext;
1952

1953
static EDealRes replacePlaceHolderColumn(SNode** pNode, void* pContext) {
147,349✔
1954
  PlaceHolderContext* pData = (PlaceHolderContext*)pContext;
147,349✔
1955
  if (QUERY_NODE_FUNCTION != nodeType((*pNode))) {
147,349✔
1956
    return DEAL_RES_CONTINUE;
121,140✔
1957
  }
1958
  SFunctionNode* pFuncNode = *(SFunctionNode**)(pNode);
26,209✔
1959
  if (!fmIsStreamPesudoColVal(pFuncNode->funcId)) {
26,209✔
1960
    return DEAL_RES_CONTINUE;
874✔
1961
  }
1962
  pData->code = fmSetStreamPseudoFuncParamVal(pFuncNode->funcId, pFuncNode->pParameterList, pData->pStreamRuntimeInfo);
25,335✔
1963
  if (pData->code != TSDB_CODE_SUCCESS) {
25,335✔
1964
    return DEAL_RES_ERROR;
×
1965
  }
1966
  SNode* pFirstParam = nodesListGetNode(pFuncNode->pParameterList, 0);
25,335✔
1967
  ((SValueNode*)pFirstParam)->translate = true;
25,335✔
1968
  SValueNode* res = NULL;
25,335✔
1969
  pData->code = nodesCloneNode(pFirstParam, (SNode**)&res);
25,335✔
1970
  if (NULL == res) {
25,335✔
1971
    return DEAL_RES_ERROR;
×
1972
  }
1973
  nodesDestroyNode(*pNode);
25,335✔
1974
  *pNode = (SNode*)res;
25,335✔
1975

1976
  return DEAL_RES_CONTINUE;
25,335✔
1977
}
1978

1979
static void extractTagColId(SOperatorNode* pOpNode, SArray* pColIdArray) {
40,040✔
1980
  SNode* pLeft = pOpNode->pLeft;
40,040✔
1981
  SNode* pRight = pOpNode->pRight;
40,040✔
1982
  SColumnNode* pColNode = nodeType(pLeft) == QUERY_NODE_COLUMN ?
40,040✔
1983
    (SColumnNode*)pLeft : (SColumnNode*)pRight;
40,040✔
1984

1985
  col_id_t colId = pColNode->colId;
40,040✔
1986
  void* _tmp = taosArrayPush(pColIdArray, &colId);
40,040✔
1987
}
40,040✔
1988

1989
static int32_t buildTagCondKey(
20,020✔
1990
  const SNode* pTagCond, char** pTagCondKey,
1991
  int32_t* tagCondKeyLen, SArray** pTagColIds) {
1992
  if (NULL == pTagCond ||
20,020✔
1993
    (nodeType(pTagCond) != QUERY_NODE_OPERATOR &&
20,020✔
1994
      nodeType(pTagCond) != QUERY_NODE_LOGIC_CONDITION)) {
20,020✔
1995
    qError("invalid parameter to extract tag filter symbol");
×
1996
    return TSDB_CODE_INTERNAL_ERROR;
×
1997
  }
1998
  int32_t code = TSDB_CODE_SUCCESS;
20,020✔
1999
  int32_t lino = 0;
20,020✔
2000
  *pTagColIds = taosArrayInit(4, sizeof(col_id_t));
20,020✔
2001

2002
  if (nodeType(pTagCond) == QUERY_NODE_OPERATOR) {
20,020✔
2003
    extractTagColId((SOperatorNode*)pTagCond, *pTagColIds);
×
2004
  } else if (nodeType(pTagCond) == QUERY_NODE_LOGIC_CONDITION) {
20,020✔
2005
    SNode* pChild = NULL;
20,020✔
2006
    FOREACH(pChild, ((SLogicConditionNode*)pTagCond)->pParameterList) {
60,060✔
2007
      extractTagColId((SOperatorNode*)pChild, *pTagColIds);
40,040✔
2008
    }
2009
  }
2010

2011
  taosArraySort(*pTagColIds, compareUint16Val);
20,020✔
2012

2013
  // encode ordered colIds into key string, separated by ','
2014
  *tagCondKeyLen =
40,040✔
2015
    (int32_t)(taosArrayGetSize(*pTagColIds) * (sizeof(col_id_t) + 1) - 1);
20,020✔
2016
  *pTagCondKey = (char*)taosMemoryCalloc(1, *tagCondKeyLen);
20,020✔
2017
  TSDB_CHECK_NULL(*pTagCondKey, code, lino, _end, terrno);
20,020✔
2018
  char* pStart = *pTagCondKey;
20,020✔
2019
  for (int32_t i = 0; i < taosArrayGetSize(*pTagColIds); ++i) {
60,060✔
2020
    col_id_t* pColId = (col_id_t*)taosArrayGet(*pTagColIds, i);
40,040✔
2021
    TSDB_CHECK_NULL(pColId, code, lino, _end, terrno);
40,040✔
2022
    memcpy(pStart, pColId, sizeof(col_id_t));
40,040✔
2023
    pStart += sizeof(col_id_t);
40,040✔
2024
    if (i != taosArrayGetSize(*pTagColIds) - 1) {
40,040✔
2025
      *pStart = ',';
20,020✔
2026
      pStart += 1;
20,020✔
2027
    }
2028
  }
2029

2030
_end:
20,020✔
2031
  if (TSDB_CODE_SUCCESS != code) {
20,020✔
2032
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
2033
    terrno = code;
×
2034
  }
2035
  return code;
20,020✔
2036
}
2037

2038
static EDealRes canOptimizeTagCondFilter(SNode* pTagCond, void* pContext) {
186,732✔
2039
  if (NULL == pTagCond) {
186,732✔
2040
    *(bool*)pContext = false;
×
2041
    return DEAL_RES_END;
×
2042
  }
2043
  if (nodeType(pTagCond) == QUERY_NODE_VALUE ||
186,732✔
2044
    nodeType(pTagCond) == QUERY_NODE_COLUMN) {
124,124✔
2045
    return DEAL_RES_CONTINUE;
103,012✔
2046
  }
2047
  if (nodeType(pTagCond) == QUERY_NODE_OPERATOR &&
83,720✔
2048
    ((SOperatorNode*)pTagCond)->opType == OP_TYPE_EQUAL) {
41,132✔
2049
    return DEAL_RES_CONTINUE;
40,040✔
2050
  }
2051
  if (nodeType(pTagCond) == QUERY_NODE_LOGIC_CONDITION &&
44,044✔
2052
    ((SLogicConditionNode*)pTagCond)->condType == LOGIC_COND_TYPE_AND) {
19,656✔
2053
    return DEAL_RES_CONTINUE;
20,020✔
2054
  }
2055
  if (nodeType(pTagCond) == QUERY_NODE_FUNCTION &&
46,956✔
2056
    fmIsStreamPesudoColVal(((SFunctionNode*)pTagCond)->funcId)) {
22,932✔
2057
    return DEAL_RES_CONTINUE;
22,932✔
2058
  }
2059
  *(bool*)pContext = false;
1,092✔
2060
  return DEAL_RES_END;
1,092✔
2061
}
2062

2063
int32_t getTableList(void* pVnode, SScanPhysiNode* pScanNode, SNode* pTagCond, SNode* pTagIndexCond,
118,650,113✔
2064
                     STableListInfo* pListInfo, uint8_t* digest, const char* idstr, SStorageAPI* pStorageAPI, void* pStreamInfo) {
2065
  int32_t code = TSDB_CODE_SUCCESS;
118,650,113✔
2066
  int32_t lino = 0;
118,650,113✔
2067
  size_t  numOfTables = 0;
118,650,113✔
2068
  bool    listAdded = false;
118,650,113✔
2069

2070
  pListInfo->idInfo.suid = pScanNode->suid;
118,682,965✔
2071
  pListInfo->idInfo.tableType = pScanNode->tableType;
118,633,550✔
2072

2073
  SArray* pUidList = taosArrayInit(8, sizeof(uint64_t));
118,627,240✔
2074
  QUERY_CHECK_NULL(pUidList, code, lino, _error, terrno);
118,643,908✔
2075

2076
  SIdxFltStatus status = SFLT_NOT_INDEX;
118,643,908✔
2077
  char*   pTagCondKey = NULL;
118,642,203✔
2078
  int32_t tagCondKeyLen;
118,646,086✔
2079
  SArray* pTagColIds = NULL;
118,655,326✔
2080
  char*   pPayload = NULL;
118,677,643✔
2081
  qTrace("getTableList called, suid:%" PRIu64
118,677,643✔
2082
    ", tagCond:%p, tagIndexCond:%p, %d %d", pScanNode->suid, pTagCond,
2083
    pTagIndexCond, pScanNode->tableType, pScanNode->virtualStableScan);
2084
  if (pScanNode->tableType != TSDB_SUPER_TABLE && !pScanNode->virtualStableScan) {
118,677,643✔
2085
    pListInfo->idInfo.uid = pScanNode->uid;
51,331,557✔
2086
    if (pStorageAPI->metaFn.isTableExisted(pVnode, pScanNode->uid)) {
51,326,154✔
2087
      void* tmp = taosArrayPush(pUidList, &pScanNode->uid);
51,280,471✔
2088
      QUERY_CHECK_NULL(tmp, code, lino, _error, terrno);
51,278,927✔
2089
    }
2090
    code = doFilterByTagCond(pListInfo, pUidList, pTagCond, pVnode, status, pStorageAPI, false, &listAdded, pStreamInfo);
51,332,950✔
2091
    QUERY_CHECK_CODE(code, lino, _end);
51,331,804✔
2092
  } else {
2093
    bool      isStream = (pStreamInfo != NULL);
67,280,010✔
2094
    bool      hasTagCond = (pTagCond != NULL);
67,280,010✔
2095
    bool      canCacheTagEqCondFilter = false;
67,280,010✔
2096
    T_MD5_CTX context = {0};
67,311,785✔
2097

2098
    qTrace("start to get table list by tag filter, suid:%" PRIu64
67,386,094✔
2099
      ",tsStableTagFilterCache:%d, tsTagFilterCache:%d", 
2100
      pScanNode->suid, tsStableTagFilterCache, tsTagFilterCache);
2101

2102
    bool acquired = false;
67,386,094✔
2103
    // first, check whether we can use stable tag filter cache
2104
    if (tsStableTagFilterCache && isStream && hasTagCond) {
67,257,051✔
2105
      canCacheTagEqCondFilter = true;
21,112✔
2106
      nodesWalkExpr(pTagCond, canOptimizeTagCondFilter,
21,112✔
2107
        (void*)&canCacheTagEqCondFilter);
2108
    }
2109
    if (canCacheTagEqCondFilter) {
67,267,819✔
2110
      qDebug("%s, stable tag filter condition can be optimized", idstr);
19,656✔
2111
      if (((SStreamRuntimeFuncInfo*)pStreamInfo)->hasPlaceHolder) {
19,656✔
2112
        SNode* tmp = NULL;
20,020✔
2113
        code = nodesCloneNode((SNode*)pTagCond, &tmp);
20,020✔
2114
        QUERY_CHECK_CODE(code, lino, _error);
20,020✔
2115

2116
        PlaceHolderContext ctx = {.code = TSDB_CODE_SUCCESS, .pStreamRuntimeInfo = (SStreamRuntimeFuncInfo*)pStreamInfo};
20,020✔
2117
        nodesRewriteExpr(&tmp, replacePlaceHolderColumn, (void*)&ctx);
20,020✔
2118
        if (TSDB_CODE_SUCCESS != ctx.code) {
20,020✔
2119
          nodesDestroyNode(tmp);
×
2120
          code = ctx.code;
×
2121
          goto _error;
×
2122
        }
2123
        code = genStableTagFilterDigest(tmp, &context);
20,020✔
2124
        nodesDestroyNode(tmp);
20,020✔
2125
      } else {
2126
        code = genStableTagFilterDigest(pTagCond, &context);
×
2127
      }
2128
      QUERY_CHECK_CODE(code, lino, _error);
20,020✔
2129

2130
      code = buildTagCondKey(
20,020✔
2131
        pTagCond, &pTagCondKey, &tagCondKeyLen, &pTagColIds);
2132
      QUERY_CHECK_CODE(code, lino, _error);
20,020✔
2133
      code = pStorageAPI->metaFn.getStableCachedTableList(
20,020✔
2134
        pVnode, pScanNode->suid, pTagCondKey, tagCondKeyLen,
20,020✔
2135
        context.digest, tListLen(context.digest), pUidList, &acquired);
2136
      QUERY_CHECK_CODE(code, lino, _error);
20,020✔
2137
    } else if (tsTagFilterCache) {
67,248,163✔
2138
      // second, try to use normal tag filter cache
2139
      qDebug("%s using normal tag filter cache", idstr);
64,723✔
2140
      if (pStreamInfo != NULL && ((SStreamRuntimeFuncInfo*)pStreamInfo)->hasPlaceHolder) {
67,126✔
2141
        SNode* tmp = NULL;
2,403✔
2142
        code = nodesCloneNode((SNode*)pTagCond, &tmp);
2,403✔
2143
        QUERY_CHECK_CODE(code, lino, _error);
2,403✔
2144

2145
        PlaceHolderContext ctx = {.code = TSDB_CODE_SUCCESS, .pStreamRuntimeInfo = (SStreamRuntimeFuncInfo*)pStreamInfo};
2,403✔
2146
        nodesRewriteExpr(&tmp, replacePlaceHolderColumn, (void*)&ctx);
2,403✔
2147
        if (TSDB_CODE_SUCCESS != ctx.code) {
2,403✔
2148
          nodesDestroyNode(tmp);
×
2149
          code = ctx.code;
×
2150
          goto _error;
×
2151
        }
2152
        code = genTagFilterDigest(tmp, &context);
2,403✔
2153
        nodesDestroyNode(tmp);
2,403✔
2154
      } else {
2155
        code = genTagFilterDigest(pTagCond, &context);
62,320✔
2156
      }
2157
      // try to retrieve the result from meta cache
2158
      QUERY_CHECK_CODE(code, lino, _error);      
64,723✔
2159
      code = pStorageAPI->metaFn.getCachedTableList(
64,723✔
2160
        pVnode, pScanNode->suid, context.digest,
64,723✔
2161
        tListLen(context.digest), pUidList, &acquired);
2162
      QUERY_CHECK_CODE(code, lino, _error);
25,303✔
2163
    }
2164
    if (acquired) {
67,211,784✔
2165
      taosArrayDestroy(pTagColIds);
58,606✔
2166
      pTagColIds = NULL;
58,606✔
2167
      
2168
      digest[0] = 1;
58,606✔
2169
      memcpy(
117,212✔
2170
        digest + 1, context.digest, tListLen(context.digest));
58,606✔
2171
      qDebug("suid:%" PRIu64 ", %s retrieve table uid list from cache,"
58,606✔
2172
        " numOfTables:%d", 
2173
        pScanNode->suid, idstr, (int32_t)taosArrayGetSize(pUidList));
2174
      goto _end;
58,606✔
2175
    } else {
2176
      qDebug("suid:%" PRIu64 
67,153,178✔
2177
        ", failed to get table uid list from cache", pScanNode->suid);
2178
    }
2179

2180
    if (!pTagCond) {  // no tag filter condition exists, let's fetch all tables of this super table
67,337,449✔
2181
      code = pStorageAPI->metaFn.getChildTableList(pVnode, pScanNode->suid, pUidList);
54,731,309✔
2182
      QUERY_CHECK_CODE(code, lino, _error);
54,695,653✔
2183
      qTrace("no tag filter, get all child tables, numOfTables:%d", (int32_t)taosArrayGetSize(pUidList));
54,695,653✔
2184
    } else {
2185
      // failed to find the result in the cache, let try to calculate the results
2186
      if (pTagIndexCond) {
12,606,140✔
2187
        void* pIndex = pStorageAPI->metaFn.getInvertIndex(pVnode);
4,375,007✔
2188

2189
        SIndexMetaArg metaArg = {.metaEx = pVnode,
4,375,073✔
2190
                                 .idx = pStorageAPI->metaFn.storeGetIndexInfo(pVnode),
4,375,007✔
2191
                                 .ivtIdx = pIndex,
2192
                                 .suid = pScanNode->uid};
4,375,007✔
2193

2194
        status = SFLT_NOT_INDEX;
4,375,007✔
2195
        code = doFilterTag(pTagIndexCond, &metaArg, pUidList, &status, &pStorageAPI->metaFilter);
4,375,007✔
2196
        if (code != 0 || status == SFLT_NOT_INDEX) {  // temporarily disable it for performance sake
4,370,118✔
2197
          qDebug("failed to get tableIds from index, suid:%" PRIu64 ", uidListSize:%d", pScanNode->uid, (int32_t)taosArrayGetSize(pUidList));
1,065,205✔
2198
        } else {
2199
          qDebug("succ to get filter result, table num: %d", (int)taosArrayGetSize(pUidList));
3,304,913✔
2200
        }
2201
      }
2202
    }
2203
    qTrace("after index filter, pTagCond:%p uidListSize:%d", pTagCond, (int32_t)taosArrayGetSize(pUidList));
67,294,728✔
2204
    code = doFilterByTagCond(pListInfo, pUidList, pTagCond, pVnode, status,
67,309,813✔
2205
      pStorageAPI, tsTagFilterCache || tsStableTagFilterCache,
67,309,813✔
2206
      &listAdded, pStreamInfo);
2207
    QUERY_CHECK_CODE(code, lino, _error);
67,301,824✔
2208

2209
    // let's add the filter results into meta-cache
2210
    numOfTables = taosArrayGetSize(pUidList);
67,300,726✔
2211

2212
    if (canCacheTagEqCondFilter) {
67,305,810✔
2213
      qInfo("%s, suid:%" PRIu64 ", add uid list to stable tag filter cache, "
10,192✔
2214
            "uidListSize:%d, origin key:%" PRIu64 ",%" PRIu64,
2215
            idstr, pScanNode->suid, (int32_t)numOfTables,
2216
            *(uint64_t*)context.digest, *(uint64_t*)(context.digest + 8));
2217

2218
      code = pStorageAPI->metaFn.putStableCachedTableList(
10,192✔
2219
        pVnode, pScanNode->suid, pTagCondKey, tagCondKeyLen,
2220
        context.digest, tListLen(context.digest),
2221
        pUidList, &pTagColIds);
2222
      QUERY_CHECK_CODE(code, lino, _end);
10,192✔
2223

2224
      digest[0] = 1;
10,192✔
2225
      memcpy(digest + 1, context.digest, tListLen(context.digest));
10,192✔
2226
    } else if (tsTagFilterCache) {
67,295,618✔
2227
      qInfo("%s, suid:%" PRIu64 ", add uid list to normal tag filter cache, "
15,945✔
2228
            "uidListSize:%d, origin key:%" PRIu64 ",%" PRIu64,
2229
            idstr, pScanNode->suid, (int32_t)numOfTables,
2230
            *(uint64_t*)context.digest, *(uint64_t*)(context.digest + 8));
2231
      size_t size = numOfTables * sizeof(uint64_t) + sizeof(int32_t);
15,945✔
2232
      pPayload = taosMemoryMalloc(size);
15,945✔
2233
      QUERY_CHECK_NULL(pPayload, code, lino, _end, terrno);
15,945✔
2234

2235
      *(int32_t*)pPayload = (int32_t)numOfTables;
15,945✔
2236
      if (numOfTables > 0) {
15,945✔
2237
        void* tmp = taosArrayGet(pUidList, 0);
12,887✔
2238
        QUERY_CHECK_NULL(tmp, code, lino, _end, terrno);
12,887✔
2239
        memcpy(pPayload + sizeof(int32_t), tmp, numOfTables * sizeof(uint64_t));
12,887✔
2240
      }
2241

2242
      code = pStorageAPI->metaFn.putCachedTableList(pVnode, pScanNode->suid,
15,945✔
2243
                                                    context.digest,
2244
                                                    tListLen(context.digest),
2245
                                                    pPayload, size, 1);
2246
      if (TSDB_CODE_SUCCESS == code) {
15,945✔
2247
        /*
2248
          data referenced by pPayload is used in lru cache,
2249
          reset pPayload to NULL to avoid being freed in _error block
2250
        */
2251
        pPayload = NULL;
15,945✔
2252
      } else {
UNCOV
2253
        if (TSDB_CODE_DUP_KEY == code) {
×
2254
          /*
2255
            another thread has already put the same key into cache,
2256
            we can just ignore this error
2257
          */
UNCOV
2258
          code = TSDB_CODE_SUCCESS;
×
2259
        }
UNCOV
2260
        QUERY_CHECK_CODE(code, lino, _end);
×
2261
      }
2262

2263

2264
      digest[0] = 1;
15,945✔
2265
      memcpy(digest + 1, context.digest, tListLen(context.digest));
15,945✔
2266
    }
2267
  }
2268

2269
_end:
118,685,904✔
2270
  if (!listAdded) {
118,672,684✔
2271
    numOfTables = taosArrayGetSize(pUidList);
107,598,585✔
2272
    for (int i = 0; i < numOfTables; i++) {
423,856,936✔
2273
      void* tmp = taosArrayGet(pUidList, i);
316,216,774✔
2274
      QUERY_CHECK_NULL(tmp, code, lino, _error, terrno);
316,255,745✔
2275
      STableKeyInfo info = {.uid = *(uint64_t*)tmp, .groupId = 0};
316,255,745✔
2276

2277
      void* p = taosArrayPush(pListInfo->pTableList, &info);
316,246,351✔
2278
      if (p == NULL) {
316,265,586✔
2279
        taosArrayDestroy(pUidList);
×
2280
        return terrno;
×
2281
      }
2282

2283
      qTrace("tagfilter get uid:%" PRIu64 ", %s", info.uid, idstr);
316,265,586✔
2284
    }
2285
  }
2286

2287
  qDebug("%s, table list with %d uids built", idstr, (int32_t)numOfTables);
118,714,261✔
2288

2289
_error:
118,732,150✔
2290
  taosArrayDestroy(pUidList);
118,739,883✔
2291
  taosArrayDestroy(pTagColIds);
118,737,482✔
2292
  taosMemFreeClear(pTagCondKey);
118,738,488✔
2293
  taosMemFreeClear(pPayload);
118,738,488✔
2294
  if (code != TSDB_CODE_SUCCESS) {
118,738,488✔
2295
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
1,098✔
2296
  }
2297
  return code;
118,732,508✔
2298
}
2299

2300
int32_t qGetTableList(int64_t suid, void* pVnode, void* node, SArray** tableList, void* pTaskInfo) {
5,138✔
2301
  int32_t        code = TSDB_CODE_SUCCESS;
5,138✔
2302
  int32_t        lino = 0;
5,138✔
2303
  SSubplan*      pSubplan = (SSubplan*)node;
5,138✔
2304
  SScanPhysiNode pNode = {0};
5,138✔
2305
  pNode.suid = suid;
5,138✔
2306
  pNode.uid = suid;
5,138✔
2307
  pNode.tableType = TSDB_SUPER_TABLE;
5,138✔
2308

2309
  STableListInfo* pTableListInfo = tableListCreate();
5,138✔
2310
  QUERY_CHECK_NULL(pTableListInfo, code, lino, _end, terrno);
5,138✔
2311
  uint8_t digest[17] = {0};
5,138✔
2312
  code = getTableList(pVnode, &pNode, pSubplan ? pSubplan->pTagCond : NULL, pSubplan ? pSubplan->pTagIndexCond : NULL,
5,138✔
2313
                      pTableListInfo, digest, "qGetTableList", &((SExecTaskInfo*)pTaskInfo)->storageAPI, NULL);
2314
  QUERY_CHECK_CODE(code, lino, _end);
5,138✔
2315
  *tableList = pTableListInfo->pTableList;
5,138✔
2316
  pTableListInfo->pTableList = NULL;
5,138✔
2317
  tableListDestroy(pTableListInfo);
5,138✔
2318

2319
_end:
5,138✔
2320
  if (code != TSDB_CODE_SUCCESS) {
5,138✔
2321
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
2322
  }
2323
  return code;
5,138✔
2324
}
2325

2326
size_t getTableTagsBufLen(const SNodeList* pGroups) {
×
2327
  size_t keyLen = 0;
×
2328

2329
  SNode* node;
2330
  FOREACH(node, pGroups) {
×
2331
    SExprNode* pExpr = (SExprNode*)node;
×
2332
    keyLen += pExpr->resType.bytes;
×
2333
  }
2334

2335
  keyLen += sizeof(int8_t) * LIST_LENGTH(pGroups);
×
2336
  return keyLen;
×
2337
}
2338

2339
int32_t getGroupIdFromTagsVal(void* pVnode, uint64_t uid, SNodeList* pGroupNode, char* keyBuf, uint64_t* pGroupId,
×
2340
                              SStorageAPI* pAPI) {
2341
  SMetaReader mr = {0};
×
2342

2343
  pAPI->metaReaderFn.initReader(&mr, pVnode, META_READER_LOCK, &pAPI->metaFn);
×
2344
  if (pAPI->metaReaderFn.getEntryGetUidCache(&mr, uid) != 0) {  // table not exist
×
2345
    pAPI->metaReaderFn.clearReader(&mr);
×
2346
    return TSDB_CODE_PAR_TABLE_NOT_EXIST;
×
2347
  }
2348

2349
  SNodeList* groupNew = NULL;
×
2350
  int32_t    code = nodesCloneList(pGroupNode, &groupNew);
×
2351
  if (TSDB_CODE_SUCCESS != code) {
×
2352
    pAPI->metaReaderFn.clearReader(&mr);
×
2353
    return code;
×
2354
  }
2355

2356
  STransTagExprCtx ctx = {.code = 0, .pReader = &mr};
×
2357
  nodesRewriteExprsPostOrder(groupNew, doTranslateTagExpr, &ctx);
×
2358
  if (TSDB_CODE_SUCCESS != ctx.code) {
×
2359
    nodesDestroyList(groupNew);
×
2360
    pAPI->metaReaderFn.clearReader(&mr);
×
2361
    return code;
×
2362
  }
2363
  char* isNull = (char*)keyBuf;
×
2364
  char* pStart = (char*)keyBuf + sizeof(int8_t) * LIST_LENGTH(pGroupNode);
×
2365

2366
  SNode*  pNode;
2367
  int32_t index = 0;
×
2368
  FOREACH(pNode, groupNew) {
×
2369
    SNode*  pNew = NULL;
×
2370
    int32_t code = scalarCalculateConstants(pNode, &pNew);
×
2371
    if (TSDB_CODE_SUCCESS == code) {
×
2372
      REPLACE_NODE(pNew);
×
2373
    } else {
2374
      nodesDestroyList(groupNew);
×
2375
      pAPI->metaReaderFn.clearReader(&mr);
×
2376
      return code;
×
2377
    }
2378

2379
    if (nodeType(pNew) != QUERY_NODE_VALUE) {
×
2380
      nodesDestroyList(groupNew);
×
2381
      pAPI->metaReaderFn.clearReader(&mr);
×
2382
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR));
×
2383
      return TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR;
×
2384
    }
2385
    SValueNode* pValue = (SValueNode*)pNew;
×
2386

2387
    if (pValue->node.resType.type == TSDB_DATA_TYPE_NULL || pValue->isNull) {
×
2388
      isNull[index++] = 1;
×
2389
      continue;
×
2390
    } else {
2391
      isNull[index++] = 0;
×
2392
      char* data = nodesGetValueFromNode(pValue);
×
2393
      if (pValue->node.resType.type == TSDB_DATA_TYPE_JSON) {
×
2394
        if (tTagIsJson(data)) {
×
2395
          terrno = TSDB_CODE_QRY_JSON_IN_GROUP_ERROR;
×
2396
          nodesDestroyList(groupNew);
×
2397
          pAPI->metaReaderFn.clearReader(&mr);
×
2398
          return terrno;
×
2399
        }
2400
        int32_t len = getJsonValueLen(data);
×
2401
        memcpy(pStart, data, len);
×
2402
        pStart += len;
×
2403
      } else if (IS_VAR_DATA_TYPE(pValue->node.resType.type)) {
×
2404
        if (IS_STR_DATA_BLOB(pValue->node.resType.type)) {
×
2405
          return TSDB_CODE_BLOB_NOT_SUPPORT_TAG;
×
2406
        }
2407
        memcpy(pStart, data, varDataTLen(data));
×
2408
        pStart += varDataTLen(data);
×
2409
      } else {
2410
        memcpy(pStart, data, pValue->node.resType.bytes);
×
2411
        pStart += pValue->node.resType.bytes;
×
2412
      }
2413
    }
2414
  }
2415

2416
  int32_t len = (int32_t)(pStart - (char*)keyBuf);
×
2417
  *pGroupId = calcGroupId(keyBuf, len);
×
2418

2419
  nodesDestroyList(groupNew);
×
2420
  pAPI->metaReaderFn.clearReader(&mr);
×
2421

2422
  return TSDB_CODE_SUCCESS;
×
2423
}
2424

2425
SArray* makeColumnArrayFromList(SNodeList* pNodeList) {
1,778,173✔
2426
  if (!pNodeList) {
1,778,173✔
2427
    return NULL;
×
2428
  }
2429

2430
  size_t  numOfCols = LIST_LENGTH(pNodeList);
1,778,173✔
2431
  SArray* pList = taosArrayInit(numOfCols, sizeof(SColumn));
1,778,637✔
2432
  if (pList == NULL) {
1,777,524✔
2433
    return NULL;
×
2434
  }
2435

2436
  for (int32_t i = 0; i < numOfCols; ++i) {
3,991,048✔
2437
    SColumnNode* pColNode = (SColumnNode*)nodesListGetNode(pNodeList, i);
2,213,550✔
2438
    if (!pColNode) {
2,214,185✔
2439
      taosArrayDestroy(pList);
×
2440
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR));
×
2441
      return NULL;
×
2442
    }
2443

2444
    // todo extract method
2445
    SColumn c = {0};
2,214,185✔
2446
    c.slotId = pColNode->slotId;
2,214,185✔
2447
    c.colId = pColNode->colId;
2,213,750✔
2448
    c.type = pColNode->node.resType.type;
2,214,385✔
2449
    c.bytes = pColNode->node.resType.bytes;
2,213,750✔
2450
    c.precision = pColNode->node.resType.precision;
2,213,750✔
2451
    c.scale = pColNode->node.resType.scale;
2,212,611✔
2452

2453
    void* tmp = taosArrayPush(pList, &c);
2,213,988✔
2454
    if (!tmp) {
2,213,988✔
2455
      taosArrayDestroy(pList);
×
2456
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
2457
      return NULL;
×
2458
    }
2459
  }
2460

2461
  return pList;
1,777,498✔
2462
}
2463

2464
int32_t extractColMatchInfo(SNodeList* pNodeList, SDataBlockDescNode* pOutputNodeList, int32_t* numOfOutputCols,
148,988,362✔
2465
                            int32_t type, SColMatchInfo* pMatchInfo) {
2466
  size_t  numOfCols = LIST_LENGTH(pNodeList);
148,988,362✔
2467
  int32_t code = TSDB_CODE_SUCCESS;
148,994,664✔
2468
  int32_t lino = 0;
148,994,664✔
2469

2470
  pMatchInfo->matchType = type;
148,994,664✔
2471

2472
  SArray* pList = taosArrayInit(numOfCols, sizeof(SColMatchItem));
149,000,081✔
2473
  if (pList == NULL) {
148,980,886✔
2474
    code = terrno;
×
2475
    return code;
×
2476
  }
2477

2478
  for (int32_t i = 0; i < numOfCols; ++i) {
899,512,916✔
2479
    STargetNode* pNode = (STargetNode*)nodesListGetNode(pNodeList, i);
750,517,010✔
2480
    QUERY_CHECK_NULL(pNode, code, lino, _end, terrno);
750,504,697✔
2481
    if (nodeType(pNode->pExpr) == QUERY_NODE_COLUMN) {
750,504,697✔
2482
      SColumnNode* pColNode = (SColumnNode*)pNode->pExpr;
746,117,337✔
2483

2484
      SColMatchItem c = {.needOutput = true};
746,136,144✔
2485
      c.colId = pColNode->colId;
746,133,483✔
2486
      c.srcSlotId = pColNode->slotId;
746,151,018✔
2487
      c.dstSlotId = pNode->slotId;
746,143,673✔
2488
      c.isPk = pColNode->isPk;
746,141,566✔
2489
      c.dataType = pColNode->node.resType;
746,154,255✔
2490
      void* tmp = taosArrayPush(pList, &c);
746,133,158✔
2491
      QUERY_CHECK_NULL(tmp, code, lino, _end, terrno);
746,133,158✔
2492
    }
2493
  }
2494

2495
  // set the output flag for each column in SColMatchInfo, according to the
2496
  *numOfOutputCols = 0;
148,995,906✔
2497
  int32_t num = LIST_LENGTH(pOutputNodeList->pSlots);
149,008,712✔
2498
  for (int32_t i = 0; i < num; ++i) {
983,121,489✔
2499
    SSlotDescNode* pNode = (SSlotDescNode*)nodesListGetNode(pOutputNodeList->pSlots, i);
834,119,976✔
2500
    QUERY_CHECK_NULL(pNode, code, lino, _end, terrno);
834,122,144✔
2501

2502
    // todo: add reserve flag check
2503
    // it is a column reserved for the arithmetic expression calculation
2504
    if (pNode->slotId >= numOfCols) {
834,122,144✔
2505
      (*numOfOutputCols) += 1;
83,598,845✔
2506
      continue;
83,599,277✔
2507
    }
2508

2509
    SColMatchItem* info = NULL;
750,554,317✔
2510
    for (int32_t j = 0; j < taosArrayGetSize(pList); ++j) {
2,147,483,647✔
2511
      info = taosArrayGet(pList, j);
2,147,483,647✔
2512
      QUERY_CHECK_NULL(info, code, lino, _end, terrno);
2,147,483,647✔
2513
      if (info->dstSlotId == pNode->slotId) {
2,147,483,647✔
2514
        break;
745,327,071✔
2515
      }
2516
    }
2517

2518
    if (pNode->output) {
13,127,998✔
2519
      (*numOfOutputCols) += 1;
744,910,544✔
2520
    } else if (info != NULL) {
5,627,862✔
2521
      // select distinct tbname from stb where tbname='abc';
2522
      info->needOutput = false;
5,632,565✔
2523
    }
2524
  }
2525

2526
  pMatchInfo->pList = pList;
149,001,513✔
2527

2528
_end:
148,982,642✔
2529
  if (code != TSDB_CODE_SUCCESS) {
148,982,642✔
2530
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
2531
  }
2532
  return code;
148,980,593✔
2533
}
2534

2535
static SResSchema createResSchema(int32_t type, int32_t bytes, int32_t slotId, int32_t scale, int32_t precision,
629,074,499✔
2536
                                  const char* name) {
2537
  SResSchema s = {0};
629,074,499✔
2538
  s.scale = scale;
629,112,054✔
2539
  s.type = type;
629,112,054✔
2540
  s.bytes = bytes;
629,112,054✔
2541
  s.slotId = slotId;
629,112,054✔
2542
  s.precision = precision;
629,112,054✔
2543
  tstrncpy(s.name, name, tListLen(s.name));
629,112,054✔
2544

2545
  return s;
629,112,054✔
2546
}
2547

2548
static SColumn* createColumn(int32_t blockId, int32_t slotId, int32_t colId, SDataType* pType, EColumnType colType) {
602,308,769✔
2549
  SColumn* pCol = taosMemoryCalloc(1, sizeof(SColumn));
602,308,769✔
2550
  if (pCol == NULL) {
602,225,951✔
2551
    return NULL;
×
2552
  }
2553

2554
  pCol->slotId = slotId;
602,225,951✔
2555
  pCol->colId = colId;
602,227,054✔
2556
  pCol->bytes = pType->bytes;
602,236,487✔
2557
  pCol->type = pType->type;
602,285,972✔
2558
  pCol->scale = pType->scale;
602,278,629✔
2559
  pCol->precision = pType->precision;
602,311,761✔
2560
  pCol->dataBlockId = blockId;
602,371,106✔
2561
  pCol->colType = colType;
602,356,452✔
2562
  return pCol;
602,393,355✔
2563
}
2564

2565
int32_t createExprFromOneNode(SExprInfo* pExp, SNode* pNode, int16_t slotId) {
629,061,731✔
2566
  int32_t code = TSDB_CODE_SUCCESS;
629,061,731✔
2567
  int32_t lino = 0;
629,061,731✔
2568
  pExp->base.numOfParams = 0;
629,061,731✔
2569
  pExp->base.pParam = NULL;
629,106,483✔
2570
  pExp->pExpr = taosMemoryCalloc(1, sizeof(tExprNode));
629,039,564✔
2571
  QUERY_CHECK_NULL(pExp->pExpr, code, lino, _end, terrno);
628,933,631✔
2572

2573
  pExp->pExpr->_function.num = 1;
628,965,683✔
2574
  pExp->pExpr->_function.functionId = -1;
628,954,549✔
2575

2576
  int32_t type = nodeType(pNode);
629,031,984✔
2577
  // it is a project query, or group by column
2578
  if (type == QUERY_NODE_COLUMN) {
629,097,994✔
2579
    pExp->pExpr->nodeType = QUERY_NODE_COLUMN;
405,328,611✔
2580
    SColumnNode* pColNode = (SColumnNode*)pNode;
405,334,615✔
2581

2582
    pExp->base.pParam = taosMemoryCalloc(1, sizeof(SFunctParam));
405,334,615✔
2583
    QUERY_CHECK_NULL(pExp->base.pParam, code, lino, _end, terrno);
405,301,824✔
2584

2585
    pExp->base.numOfParams = 1;
405,285,836✔
2586

2587
    SDataType* pType = &pColNode->node.resType;
405,291,094✔
2588
    pExp->base.resSchema =
2589
        createResSchema(pType->type, pType->bytes, slotId, pType->scale, pType->precision, pColNode->colName);
405,263,897✔
2590

2591
    pExp->base.pParam[0].pCol =
810,617,671✔
2592
        createColumn(pColNode->dataBlockId, pColNode->slotId, pColNode->colId, pType, pColNode->colType);
810,544,004✔
2593
    QUERY_CHECK_NULL(pExp->base.pParam[0].pCol, code, lino, _end, terrno);
405,217,372✔
2594

2595
    pExp->base.pParam[0].type = FUNC_PARAM_TYPE_COLUMN;
405,287,893✔
2596
  } else if (type == QUERY_NODE_VALUE) {
223,769,383✔
2597
    pExp->pExpr->nodeType = QUERY_NODE_VALUE;
8,016,678✔
2598
    SValueNode* pValNode = (SValueNode*)pNode;
8,018,044✔
2599

2600
    pExp->base.pParam = taosMemoryCalloc(1, sizeof(SFunctParam));
8,018,044✔
2601
    QUERY_CHECK_NULL(pExp->base.pParam, code, lino, _end, terrno);
8,016,459✔
2602

2603
    pExp->base.numOfParams = 1;
8,014,844✔
2604

2605
    SDataType* pType = &pValNode->node.resType;
8,016,254✔
2606
    pExp->base.resSchema =
2607
        createResSchema(pType->type, pType->bytes, slotId, pType->scale, pType->precision, pValNode->node.aliasName);
8,016,425✔
2608
    pExp->base.pParam[0].type = FUNC_PARAM_TYPE_VALUE;
8,018,727✔
2609
    code = nodesValueNodeToVariant(pValNode, &pExp->base.pParam[0].param);
8,016,731✔
2610
    QUERY_CHECK_CODE(code, lino, _end);
8,015,767✔
2611
  } else if (type == QUERY_NODE_FUNCTION) {
215,752,705✔
2612
    pExp->pExpr->nodeType = QUERY_NODE_FUNCTION;
203,794,928✔
2613
    SFunctionNode* pFuncNode = (SFunctionNode*)pNode;
203,816,471✔
2614

2615
    SDataType* pType = &pFuncNode->node.resType;
203,816,471✔
2616
    pExp->base.resSchema =
2617
        createResSchema(pType->type, pType->bytes, slotId, pType->scale, pType->precision, pFuncNode->node.aliasName);
203,804,192✔
2618
    tExprNode* pExprNode = pExp->pExpr;
203,798,453✔
2619

2620
    pExprNode->_function.functionId = pFuncNode->funcId;
203,802,645✔
2621
    pExprNode->_function.pFunctNode = pFuncNode;
203,808,722✔
2622
    pExprNode->_function.functionType = pFuncNode->funcType;
203,800,992✔
2623

2624
    tstrncpy(pExprNode->_function.functionName, pFuncNode->functionName, tListLen(pExprNode->_function.functionName));
203,777,804✔
2625

2626
    pExp->base.pParamList = pFuncNode->pParameterList;
203,781,520✔
2627
#if 1
2628
    // todo refactor: add the parameter for tbname function
2629
    const char* name = "tbname";
203,805,531✔
2630
    int32_t     len = strlen(name);
203,805,531✔
2631

2632
    if (!pFuncNode->pParameterList && (memcmp(pExprNode->_function.functionName, name, len) == 0) &&
203,805,531✔
2633
        pExprNode->_function.functionName[len] == 0) {
9,058,120✔
2634
      pFuncNode->pParameterList = NULL;
9,058,308✔
2635
      int32_t     code = nodesMakeList(&pFuncNode->pParameterList);
9,057,949✔
2636
      SValueNode* res = NULL;
9,060,656✔
2637
      if (TSDB_CODE_SUCCESS == code) {
9,060,656✔
2638
        code = nodesMakeNode(QUERY_NODE_VALUE, (SNode**)&res);
9,057,014✔
2639
      }
2640
      QUERY_CHECK_CODE(code, lino, _end);
9,058,215✔
2641
      res->node.resType = (SDataType){.bytes = sizeof(int64_t), .type = TSDB_DATA_TYPE_BIGINT};
9,058,215✔
2642
      code = nodesListAppend(pFuncNode->pParameterList, (SNode*)res);
9,056,146✔
2643
      if (code != TSDB_CODE_SUCCESS) {
9,059,714✔
2644
        nodesDestroyNode((SNode*)res);
×
2645
        res = NULL;
×
2646
      }
2647
      QUERY_CHECK_CODE(code, lino, _end);
9,059,714✔
2648
    }
2649
#endif
2650

2651
    int32_t numOfParam = LIST_LENGTH(pFuncNode->pParameterList);
203,811,617✔
2652

2653
    pExp->base.pParam = taosMemoryCalloc(numOfParam, sizeof(SFunctParam));
203,813,116✔
2654
    QUERY_CHECK_NULL(pExp->base.pParam, code, lino, _end, terrno);
203,761,160✔
2655
    pExp->base.numOfParams = numOfParam;
203,769,151✔
2656

2657
    for (int32_t j = 0; j < numOfParam && TSDB_CODE_SUCCESS == code; ++j) {
524,999,164✔
2658
      SNode* p1 = nodesListGetNode(pFuncNode->pParameterList, j);
321,210,454✔
2659
      QUERY_CHECK_NULL(p1, code, lino, _end, terrno);
321,237,042✔
2660
      if (p1->type == QUERY_NODE_COLUMN) {
321,237,042✔
2661
        SColumnNode* pcn = (SColumnNode*)p1;
197,051,801✔
2662

2663
        pExp->base.pParam[j].type = FUNC_PARAM_TYPE_COLUMN;
197,051,801✔
2664
        pExp->base.pParam[j].pCol =
394,076,423✔
2665
            createColumn(pcn->dataBlockId, pcn->slotId, pcn->colId, &pcn->node.resType, pcn->colType);
394,061,403✔
2666
        QUERY_CHECK_NULL(pExp->base.pParam[j].pCol, code, lino, _end, terrno);
197,042,958✔
2667
      } else if (p1->type == QUERY_NODE_VALUE) {
124,177,705✔
2668
        SValueNode* pvn = (SValueNode*)p1;
63,484,395✔
2669
        pExp->base.pParam[j].type = FUNC_PARAM_TYPE_VALUE;
63,484,395✔
2670
        code = nodesValueNodeToVariant(pvn, &pExp->base.pParam[j].param);
63,483,169✔
2671
        QUERY_CHECK_CODE(code, lino, _end);
63,443,281✔
2672
      }
2673
    }
2674
    pExp->pExpr->_function.bindExprID = ((SExprNode*)pNode)->bindExprID;
203,788,710✔
2675
  } else if (type == QUERY_NODE_OPERATOR) {
11,957,777✔
2676
    pExp->pExpr->nodeType = QUERY_NODE_OPERATOR;
10,922,602✔
2677
    SOperatorNode* pOpNode = (SOperatorNode*)pNode;
10,921,837✔
2678

2679
    pExp->base.pParam = taosMemoryCalloc(1, sizeof(SFunctParam));
10,921,837✔
2680
    QUERY_CHECK_NULL(pExp->base.pParam, code, lino, _end, terrno);
10,921,417✔
2681
    pExp->base.numOfParams = 1;
10,919,423✔
2682

2683
    SDataType* pType = &pOpNode->node.resType;
10,920,032✔
2684
    pExp->base.resSchema =
2685
        createResSchema(pType->type, pType->bytes, slotId, pType->scale, pType->precision, pOpNode->node.aliasName);
10,919,036✔
2686
    pExp->pExpr->_optrRoot.pRootNode = pNode;
10,920,738✔
2687
  } else if (type == QUERY_NODE_CASE_WHEN) {
1,035,234✔
2688
    pExp->pExpr->nodeType = QUERY_NODE_OPERATOR;
1,034,467✔
2689
    SCaseWhenNode* pCaseNode = (SCaseWhenNode*)pNode;
1,034,467✔
2690

2691
    pExp->base.pParam = taosMemoryCalloc(1, sizeof(SFunctParam));
1,034,467✔
2692
    QUERY_CHECK_NULL(pExp->base.pParam, code, lino, _end, terrno);
1,034,467✔
2693
    pExp->base.numOfParams = 1;
1,034,467✔
2694

2695
    SDataType* pType = &pCaseNode->node.resType;
1,034,467✔
2696
    pExp->base.resSchema =
2697
        createResSchema(pType->type, pType->bytes, slotId, pType->scale, pType->precision, pCaseNode->node.aliasName);
1,034,467✔
2698
    pExp->pExpr->_optrRoot.pRootNode = pNode;
1,034,467✔
2699
  } else if (type == QUERY_NODE_LOGIC_CONDITION) {
971✔
2700
    pExp->pExpr->nodeType = QUERY_NODE_OPERATOR;
1,140✔
2701
    SLogicConditionNode* pCond = (SLogicConditionNode*)pNode;
1,140✔
2702
    pExp->base.pParam = taosMemoryCalloc(1, sizeof(SFunctParam));
1,140✔
2703
    QUERY_CHECK_NULL(pExp->base.pParam, code, lino, _end, terrno);
1,140✔
2704
    pExp->base.numOfParams = 1;
1,140✔
2705
    SDataType* pType = &pCond->node.resType;
1,140✔
2706
    pExp->base.resSchema =
2707
        createResSchema(pType->type, pType->bytes, slotId, pType->scale, pType->precision, pCond->node.aliasName);
1,140✔
2708
    pExp->pExpr->_optrRoot.pRootNode = pNode;
1,140✔
2709
  } else {
2710
    code = TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR;
135✔
2711
    QUERY_CHECK_CODE(code, lino, _end);
135✔
2712
  }
2713
  pExp->pExpr->relatedTo = ((SExprNode*)pNode)->relatedTo;
629,075,266✔
2714
_end:
629,090,868✔
2715
  if (code != TSDB_CODE_SUCCESS) {
629,090,868✔
2716
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
2717
  }
2718
  return code;
628,950,854✔
2719
}
2720

2721
int32_t createExprFromTargetNode(SExprInfo* pExp, STargetNode* pTargetNode) {
629,011,205✔
2722
  return createExprFromOneNode(pExp, pTargetNode->pExpr, pTargetNode->slotId);
629,011,205✔
2723
}
2724

2725
SExprInfo* createExpr(SNodeList* pNodeList, int32_t* numOfExprs) {
×
2726
  *numOfExprs = LIST_LENGTH(pNodeList);
×
2727
  SExprInfo* pExprs = taosMemoryCalloc(*numOfExprs, sizeof(SExprInfo));
×
2728
  if (!pExprs) {
×
2729
    return NULL;
×
2730
  }
2731

2732
  for (int32_t i = 0; i < (*numOfExprs); ++i) {
×
2733
    SExprInfo* pExp = &pExprs[i];
×
2734
    int32_t    code = createExprFromOneNode(pExp, nodesListGetNode(pNodeList, i), i + UD_TAG_COLUMN_INDEX);
×
2735
    if (code != TSDB_CODE_SUCCESS) {
×
2736
      taosMemoryFreeClear(pExprs);
×
2737
      terrno = code;
×
2738
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(code));
×
2739
      return NULL;
×
2740
    }
2741
  }
2742

2743
  return pExprs;
×
2744
}
2745

2746
int32_t createExprInfo(SNodeList* pNodeList, SNodeList* pGroupKeys, SExprInfo** pExprInfo, int32_t* numOfExprs) {
207,117,448✔
2747
  QRY_PARAM_CHECK(pExprInfo);
207,117,448✔
2748

2749
  int32_t code = 0;
207,147,658✔
2750
  int32_t numOfFuncs = LIST_LENGTH(pNodeList);
207,147,658✔
2751
  int32_t numOfGroupKeys = 0;
207,129,422✔
2752
  if (pGroupKeys != NULL) {
207,129,422✔
2753
    numOfGroupKeys = LIST_LENGTH(pGroupKeys);
18,919,960✔
2754
  }
2755

2756
  *numOfExprs = numOfFuncs + numOfGroupKeys;
207,130,805✔
2757
  if (*numOfExprs == 0) {
207,127,347✔
2758
    return code;
19,523,862✔
2759
  }
2760

2761
  SExprInfo* pExprs = taosMemoryCalloc(*numOfExprs, sizeof(SExprInfo));
187,616,758✔
2762
  if (pExprs == NULL) {
187,538,603✔
2763
    return terrno;
×
2764
  }
2765

2766
  for (int32_t i = 0; i < (*numOfExprs); ++i) {
816,414,863✔
2767
    STargetNode* pTargetNode = NULL;
628,896,500✔
2768
    if (i < numOfFuncs) {
628,896,500✔
2769
      pTargetNode = (STargetNode*)nodesListGetNode(pNodeList, i);
605,695,093✔
2770
    } else {
2771
      pTargetNode = (STargetNode*)nodesListGetNode(pGroupKeys, i - numOfFuncs);
23,201,407✔
2772
    }
2773
    if (!pTargetNode) {
628,932,830✔
2774
      destroyExprInfo(pExprs, *numOfExprs);
×
2775
      taosMemoryFreeClear(pExprs);
×
2776
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
2777
      return terrno;
×
2778
    }
2779

2780
    SExprInfo* pExp = &pExprs[i];
628,932,830✔
2781
    code = createExprFromTargetNode(pExp, pTargetNode);
628,938,026✔
2782
    if (code != TSDB_CODE_SUCCESS) {
628,876,260✔
2783
      destroyExprInfo(pExprs, *numOfExprs);
×
2784
      taosMemoryFreeClear(pExprs);
×
2785
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(code));
×
2786
      return code;
×
2787
    }
2788
  }
2789

2790
  *pExprInfo = pExprs;
187,591,135✔
2791
  return code;
187,602,598✔
2792
}
2793

2794
static void deleteSubsidiareCtx(void* pData) {
×
2795
  SSubsidiaryResInfo* pCtx = (SSubsidiaryResInfo*)pData;
×
2796
  if (pCtx->pCtx) {
×
2797
    taosMemoryFreeClear(pCtx->pCtx);
×
2798
  }
2799
}
×
2800

2801
// set the output buffer for the selectivity + tag query
2802
static int32_t setSelectValueColumnInfo(SqlFunctionCtx* pCtx, int32_t numOfOutput) {
197,347,862✔
2803
  int32_t num = 0;
197,347,862✔
2804
  int32_t code = TSDB_CODE_SUCCESS;
197,347,862✔
2805
  int32_t lino = 0;
197,347,862✔
2806

2807
  SArray* pValCtxArray = NULL;
197,347,862✔
2808
  for (int32_t i = numOfOutput - 1; i > 0; --i) {  // select Func is at the end of the list
634,342,504✔
2809
    int32_t funcIdx = pCtx[i].pExpr->pExpr->_function.bindExprID;
437,017,476✔
2810
    if (funcIdx > 0) {
437,045,273✔
2811
      if (pValCtxArray == NULL) {
1,838,130✔
2812
        // the end of the list is the select function of biggest index
2813
        pValCtxArray = taosArrayInit_s(sizeof(SSubsidiaryResInfo*), funcIdx);
1,318,275✔
2814
        if (pValCtxArray == NULL) {
1,317,269✔
2815
          return terrno;
×
2816
        }
2817
      }
2818
      if (funcIdx > pValCtxArray->size) {
1,837,124✔
2819
        qError("funcIdx:%d is out of range", funcIdx);
×
2820
        taosArrayDestroyP(pValCtxArray, deleteSubsidiareCtx);
×
2821
        return TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR;
×
2822
      }
2823
      SSubsidiaryResInfo* pSubsidiary = &pCtx[i].subsidiaries;
1,837,124✔
2824
      pSubsidiary->pCtx = taosMemoryCalloc(numOfOutput, POINTER_BYTES);
1,840,645✔
2825
      if (pSubsidiary->pCtx == NULL) {
1,839,639✔
2826
        taosArrayDestroyP(pValCtxArray, deleteSubsidiareCtx);
×
2827
        return terrno;
×
2828
      }
2829
      pSubsidiary->num = 0;
1,838,633✔
2830
      taosArraySet(pValCtxArray, funcIdx - 1, &pSubsidiary);
1,839,136✔
2831
    }
2832
  }
2833

2834
  SqlFunctionCtx*  p = NULL;
197,325,028✔
2835
  SqlFunctionCtx** pValCtx = NULL;
197,325,028✔
2836
  if (pValCtxArray == NULL) {
197,325,028✔
2837
    pValCtx = taosMemoryCalloc(numOfOutput, POINTER_BYTES);
196,010,060✔
2838
    if (pValCtx == NULL) {
196,073,070✔
2839
      QUERY_CHECK_CODE(terrno, lino, _end);
×
2840
    }
2841
  }
2842

2843
  for (int32_t i = 0; i < numOfOutput; ++i) {
818,934,393✔
2844
    const char* pName = pCtx[i].pExpr->pExpr->_function.functionName;
621,521,170✔
2845
    if ((strcmp(pName, "_select_value") == 0)) {
621,488,612✔
2846
      if (pValCtxArray == NULL) {
5,894,827✔
2847
        pValCtx[num++] = &pCtx[i];
3,313,226✔
2848
      } else {
2849
        int32_t bindFuncIndex = pCtx[i].pExpr->pExpr->relatedTo;  // start from index 1;
2,581,750✔
2850
        if (bindFuncIndex > 0) {                                  // 0 is default index related to the select function
2,581,276✔
2851
          bindFuncIndex -= 1;
2,522,425✔
2852
        }
2853
        SSubsidiaryResInfo** pSubsidiary = taosArrayGet(pValCtxArray, bindFuncIndex);
2,581,276✔
2854
        if (pSubsidiary == NULL) {
2,581,779✔
2855
          QUERY_CHECK_CODE(TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR, lino, _end);
×
2856
        }
2857
        (*pSubsidiary)->pCtx[(*pSubsidiary)->num] = &pCtx[i];
2,581,779✔
2858
        (*pSubsidiary)->num++;
2,581,276✔
2859
      }
2860
    } else if (fmIsSelectFunc(pCtx[i].functionId)) {
615,593,785✔
2861
      if (pValCtxArray == NULL) {
51,862,893✔
2862
        p = &pCtx[i];
49,625,912✔
2863
      }
2864
    }
2865
  }
2866

2867
  if (p != NULL) {
197,413,223✔
2868
    p->subsidiaries.pCtx = pValCtx;
20,177,090✔
2869
    p->subsidiaries.num = num;
20,176,682✔
2870
  } else {
2871
    taosMemoryFreeClear(pValCtx);
177,236,133✔
2872
  }
2873

2874
_end:
1,362,092✔
2875
  if (code != TSDB_CODE_SUCCESS) {
197,375,245✔
2876
    taosArrayDestroyP(pValCtxArray, deleteSubsidiareCtx);
×
2877
    taosMemoryFreeClear(pValCtx);
×
2878
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
2879
  } else {
2880
    taosArrayDestroy(pValCtxArray);
197,375,245✔
2881
  }
2882
  return code;
197,352,953✔
2883
}
2884

2885
SqlFunctionCtx* createSqlFunctionCtx(SExprInfo* pExprInfo, int32_t numOfOutput, int32_t** rowEntryInfoOffset,
197,387,385✔
2886
                                     SFunctionStateStore* pStore) {
2887
  int32_t         code = TSDB_CODE_SUCCESS;
197,387,385✔
2888
  int32_t         lino = 0;
197,387,385✔
2889
  SqlFunctionCtx* pFuncCtx = (SqlFunctionCtx*)taosMemoryCalloc(numOfOutput, sizeof(SqlFunctionCtx));
197,387,385✔
2890
  if (pFuncCtx == NULL) {
197,341,348✔
2891
    return NULL;
×
2892
  }
2893

2894
  *rowEntryInfoOffset = taosMemoryCalloc(numOfOutput, sizeof(int32_t));
197,341,348✔
2895
  if (*rowEntryInfoOffset == 0) {
197,437,131✔
2896
    taosMemoryFreeClear(pFuncCtx);
×
2897
    return NULL;
×
2898
  }
2899

2900
  for (int32_t i = 0; i < numOfOutput; ++i) {
819,023,210✔
2901
    SExprInfo* pExpr = &pExprInfo[i];
621,622,522✔
2902

2903
    SExprBasicInfo* pFunct = &pExpr->base;
621,574,820✔
2904
    SqlFunctionCtx* pCtx = &pFuncCtx[i];
621,621,329✔
2905

2906
    pCtx->functionId = -1;
621,630,223✔
2907
    pCtx->pExpr = pExpr;
621,611,243✔
2908

2909
    if (pExpr->pExpr->nodeType == QUERY_NODE_FUNCTION) {
621,621,100✔
2910
      SFuncExecEnv env = {0};
202,447,525✔
2911
      pCtx->functionId = pExpr->pExpr->_function.pFunctNode->funcId;
202,442,773✔
2912
      pCtx->isPseudoFunc = fmIsWindowPseudoColumnFunc(pCtx->functionId) || fmIsPlaceHolderFunc(pCtx->functionId);
202,417,574✔
2913
      pCtx->isNotNullFunc = fmIsNotNullOutputFunc(pCtx->functionId);
202,461,582✔
2914

2915
      bool isUdaf = fmIsUserDefinedFunc(pCtx->functionId);
202,434,982✔
2916
      if (fmIsAggFunc(pCtx->functionId) || fmIsIndefiniteRowsFunc(pCtx->functionId)) {
323,412,221✔
2917
        if (!isUdaf) {
121,011,059✔
2918
          code = fmGetFuncExecFuncs(pCtx->functionId, &pCtx->fpSet);
120,969,166✔
2919
          QUERY_CHECK_CODE(code, lino, _end);
120,949,356✔
2920
        } else {
2921
          char* udfName = pExpr->pExpr->_function.pFunctNode->functionName;
41,893✔
2922
          pCtx->udfName = taosStrdup(udfName);
41,893✔
2923
          QUERY_CHECK_NULL(pCtx->udfName, code, lino, _end, terrno);
41,893✔
2924

2925
          code = fmGetUdafExecFuncs(pCtx->functionId, &pCtx->fpSet);
41,893✔
2926
          QUERY_CHECK_CODE(code, lino, _end);
41,893✔
2927
        }
2928
        bool tmp = pCtx->fpSet.getEnv(pExpr->pExpr->_function.pFunctNode, &env);
120,991,249✔
2929
        if (!tmp) {
121,000,655✔
2930
          code = terrno;
×
UNCOV
2931
          QUERY_CHECK_CODE(code, lino, _end);
×
2932
        }
2933
      } else {
2934
        if (fmIsPlaceHolderFunc(pCtx->functionId)) {
81,422,309✔
2935
          code = fmGetStreamPesudoFuncEnv(pCtx->functionId, pExpr->base.pParamList, &env);
7,696,651✔
2936
          QUERY_CHECK_CODE(code, lino, _end);
7,696,622✔
2937
        }      
2938
        
2939
        code = fmGetScalarFuncExecFuncs(pCtx->functionId, &pCtx->sfp);
81,426,851✔
2940
        if (code != TSDB_CODE_SUCCESS && isUdaf) {
81,440,321✔
2941
          code = TSDB_CODE_SUCCESS;
25,956✔
2942
        }
2943
        QUERY_CHECK_CODE(code, lino, _end);
81,440,321✔
2944

2945
        if (pCtx->sfp.getEnv != NULL) {
81,440,321✔
2946
          bool tmp = pCtx->sfp.getEnv(pExpr->pExpr->_function.pFunctNode, &env);
15,661,553✔
2947
          if (!tmp) {
15,661,410✔
2948
            code = terrno;
×
2949
            QUERY_CHECK_CODE(code, lino, _end);
×
2950
          }
2951
        }
2952
      }
2953
      pCtx->resDataInfo.interBufSize = env.calcMemSize;
202,436,034✔
2954
    } else if (pExpr->pExpr->nodeType == QUERY_NODE_COLUMN || pExpr->pExpr->nodeType == QUERY_NODE_OPERATOR ||
419,145,886✔
2955
               pExpr->pExpr->nodeType == QUERY_NODE_VALUE) {
8,016,375✔
2956
      // for simple column, the result buffer needs to hold at least one element.
2957
      pCtx->resDataInfo.interBufSize = pFunct->resSchema.bytes;
419,165,565✔
2958
    }
2959

2960
    pCtx->input.numOfInputCols = pFunct->numOfParams;
621,658,121✔
2961
    pCtx->input.pData = taosMemoryCalloc(pFunct->numOfParams, POINTER_BYTES);
621,614,368✔
2962
    QUERY_CHECK_NULL(pCtx->input.pData, code, lino, _end, terrno);
621,627,980✔
2963
    pCtx->input.pColumnDataAgg = taosMemoryCalloc(pFunct->numOfParams, POINTER_BYTES);
621,603,340✔
2964
    QUERY_CHECK_NULL(pCtx->input.pColumnDataAgg, code, lino, _end, terrno);
621,631,962✔
2965

2966
    pCtx->pTsOutput = NULL;
621,610,016✔
2967
    pCtx->resDataInfo.bytes = pFunct->resSchema.bytes;
621,627,518✔
2968
    pCtx->resDataInfo.type = pFunct->resSchema.type;
621,643,990✔
2969
    pCtx->order = TSDB_ORDER_ASC;
621,625,074✔
2970
    pCtx->start.key = INT64_MIN;
621,632,610✔
2971
    pCtx->end.key = INT64_MIN;
621,595,474✔
2972
    pCtx->numOfParams = pExpr->base.numOfParams;
621,646,643✔
2973
    pCtx->param = pFunct->pParam;
621,615,683✔
2974
    pCtx->saveHandle.currentPage = -1;
621,654,969✔
2975
    pCtx->pStore = pStore;
621,645,835✔
2976
    pCtx->hasWindowOrGroup = false;
621,626,686✔
2977
    pCtx->needCleanup = false;
621,600,945✔
2978
  }
2979

2980
  for (int32_t i = 1; i < numOfOutput; ++i) {
634,407,321✔
2981
    (*rowEntryInfoOffset)[i] = (int32_t)((*rowEntryInfoOffset)[i - 1] + sizeof(SResultRowEntryInfo) +
874,105,359✔
2982
                                         pFuncCtx[i - 1].resDataInfo.interBufSize);
437,083,041✔
2983
  }
2984

2985
  code = setSelectValueColumnInfo(pFuncCtx, numOfOutput);
197,350,888✔
2986
  QUERY_CHECK_CODE(code, lino, _end);
197,357,747✔
2987

2988
_end:
197,357,747✔
2989
  if (code != TSDB_CODE_SUCCESS) {
197,323,137✔
2990
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
2991
    for (int32_t i = 0; i < numOfOutput; ++i) {
×
2992
      taosMemoryFree(pFuncCtx[i].input.pData);
×
2993
      taosMemoryFree(pFuncCtx[i].input.pColumnDataAgg);
×
2994
    }
2995
    taosMemoryFreeClear(*rowEntryInfoOffset);
×
2996
    taosMemoryFreeClear(pFuncCtx);
×
2997

2998
    terrno = code;
×
2999
    return NULL;
×
3000
  }
3001
  return pFuncCtx;
197,323,137✔
3002
}
3003

3004
// NOTE: sources columns are more than the destination SSDatablock columns.
3005
// doFilter in table scan needs every column even its output is false
3006
int32_t relocateColumnData(SSDataBlock* pBlock, const SArray* pColMatchInfo, SArray* pCols, bool outputEveryColumn) {
9,725,873✔
3007
  int32_t code = TSDB_CODE_SUCCESS;
9,725,873✔
3008
  size_t  numOfSrcCols = taosArrayGetSize(pCols);
9,725,873✔
3009

3010
  int32_t i = 0, j = 0;
9,726,374✔
3011
  while (i < numOfSrcCols && j < taosArrayGetSize(pColMatchInfo)) {
89,572,161✔
3012
    SColumnInfoData* p = taosArrayGet(pCols, i);
79,846,789✔
3013
    if (!p) {
79,847,404✔
3014
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
3015
      return terrno;
×
3016
    }
3017
    SColMatchItem* pmInfo = taosArrayGet(pColMatchInfo, j);
79,847,404✔
3018
    if (!pmInfo) {
79,845,197✔
3019
      return terrno;
×
3020
    }
3021

3022
    if (p->info.colId == pmInfo->colId) {
79,845,197✔
3023
      SColumnInfoData* pDst = taosArrayGet(pBlock->pDataBlock, pmInfo->dstSlotId);
71,133,416✔
3024
      if (!pDst) {
71,132,915✔
3025
        return terrno;
×
3026
      }
3027
      code = colDataAssign(pDst, p, pBlock->info.rows, &pBlock->info);
71,132,915✔
3028
      if (code != TSDB_CODE_SUCCESS) {
71,131,762✔
3029
        qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(code));
×
3030
        return code;
×
3031
      }
3032
      i++;
71,131,762✔
3033
      j++;
71,131,762✔
3034
    } else if (p->info.colId < pmInfo->colId) {
8,714,025✔
3035
      i++;
8,714,025✔
3036
    } else {
3037
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR));
×
3038
      return TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR;
×
3039
    }
3040
  }
3041
  return code;
9,726,374✔
3042
}
3043

3044
SInterval extractIntervalInfo(const STableScanPhysiNode* pTableScanNode) {
96,824,981✔
3045
  SInterval interval = {
193,627,546✔
3046
      .interval = pTableScanNode->interval,
96,850,578✔
3047
      .sliding = pTableScanNode->sliding,
96,839,594✔
3048
      .intervalUnit = pTableScanNode->intervalUnit,
96,840,040✔
3049
      .slidingUnit = pTableScanNode->slidingUnit,
96,837,483✔
3050
      .offset = pTableScanNode->offset,
96,834,928✔
3051
      .precision = pTableScanNode->scan.node.pOutputDataBlockDesc->precision,
96,852,328✔
3052
      .timeRange = pTableScanNode->scanRange,
3053
  };
3054
  calcIntervalAutoOffset(&interval);
96,840,404✔
3055

3056
  return interval;
96,823,394✔
3057
}
3058

3059
SColumn extractColumnFromColumnNode(SColumnNode* pColNode) {
25,345,678✔
3060
  SColumn c = {0};
25,345,678✔
3061

3062
  c.slotId = pColNode->slotId;
25,345,678✔
3063
  c.colId = pColNode->colId;
25,352,229✔
3064
  c.type = pColNode->node.resType.type;
25,350,191✔
3065
  c.bytes = pColNode->node.resType.bytes;
25,352,595✔
3066
  c.scale = pColNode->node.resType.scale;
25,344,803✔
3067
  c.precision = pColNode->node.resType.precision;
25,346,675✔
3068
  return c;
25,335,382✔
3069
}
3070

3071

3072
/**
3073
 * @brief Determine the actual time range for reading data based on the RANGE clause and the WHERE conditions.
3074
 * @param[in] cond The range specified by WHERE condition.
3075
 * @param[in] range The range specified by RANGE clause.
3076
 * @param[out] twindow The range to be read in DESC order, and only one record is needed.
3077
 * @param[out] extTwindow The external range to read for only one record, which is used for FILL clause.
3078
 * @note `cond` and `twindow` may be the same address.
3079
 */
3080
static int32_t getQueryExtWindow(const STimeWindow* cond, const STimeWindow* range, STimeWindow* twindow,
1,625,418✔
3081
                                 STimeWindow* extTwindows) {
3082
  int32_t     code = TSDB_CODE_SUCCESS;
1,625,418✔
3083
  int32_t     lino = 0;
1,625,418✔
3084
  STimeWindow tempWindow;
3085

3086
  if (cond->skey > cond->ekey || range->skey > range->ekey) {
1,625,418✔
3087
    *twindow = extTwindows[0] = extTwindows[1] = TSWINDOW_DESC_INITIALIZER;
2,940✔
3088
    return code;
2,940✔
3089
  }
3090

3091
  if (range->ekey < cond->skey) {
1,622,478✔
3092
    extTwindows[1] = *cond;
244,571✔
3093
    *twindow = extTwindows[0] = TSWINDOW_DESC_INITIALIZER;
244,571✔
3094
    return code;
244,571✔
3095
  }
3096

3097
  if (cond->ekey < range->skey) {
1,377,907✔
3098
    extTwindows[0] = *cond;
190,602✔
3099
    *twindow = extTwindows[1] = TSWINDOW_DESC_INITIALIZER;
190,602✔
3100
    return code;
190,602✔
3101
  }
3102

3103
  // Only scan data in the time range intersecion.
3104
  extTwindows[0] = extTwindows[1] = *cond;
1,187,305✔
3105
  twindow->skey = TMAX(cond->skey, range->skey);
1,187,305✔
3106
  twindow->ekey = TMIN(cond->ekey, range->ekey);
1,187,305✔
3107
  extTwindows[0].ekey = twindow->skey - 1;
1,187,305✔
3108
  extTwindows[1].skey = twindow->ekey + 1;
1,187,305✔
3109

3110
  return code;
1,187,305✔
3111
}
3112

3113
int32_t initQueryTableDataCond(SQueryTableDataCond* pCond, const STableScanPhysiNode* pTableScanNode,
118,245,401✔
3114
                               const SReadHandle* readHandle, bool applyExtWin) {
3115
  int32_t code = 0;                             
118,245,401✔
3116
  pCond->order = pTableScanNode->scanSeq[0] > 0 ? TSDB_ORDER_ASC : TSDB_ORDER_DESC;
118,245,401✔
3117
  pCond->numOfCols = LIST_LENGTH(pTableScanNode->scan.pScanCols);
118,245,314✔
3118

3119
  pCond->colList = taosMemoryCalloc(pCond->numOfCols, sizeof(SColumnInfo));
118,239,556✔
3120
  if (!pCond->colList) {
118,232,577✔
3121
    return terrno;
×
3122
  }
3123
  pCond->pSlotList = taosMemoryMalloc(sizeof(int32_t) * pCond->numOfCols);
118,222,732✔
3124
  if (pCond->pSlotList == NULL) {
118,220,655✔
3125
    taosMemoryFreeClear(pCond->colList);
×
3126
    return terrno;
×
3127
  }
3128

3129
  // TODO: get it from stable scan node
3130
  pCond->twindows = pTableScanNode->scanRange;
118,195,214✔
3131
  pCond->suid = pTableScanNode->scan.suid;
118,256,657✔
3132
  pCond->type = TIMEWINDOW_RANGE_CONTAINED;
118,221,409✔
3133
  pCond->startVersion = -1;
118,198,995✔
3134
  pCond->endVersion = -1;
118,213,420✔
3135
  pCond->skipRollup = readHandle->skipRollup;
118,209,313✔
3136
  if (readHandle->winRangeValid) {
118,189,814✔
3137
    pCond->twindows = readHandle->winRange;
339,603✔
3138
  }
3139
  pCond->cacheSttStatis = readHandle->cacheSttStatis;
118,241,377✔
3140
  // allowed read stt file optimization mode
3141
  pCond->notLoadData = (pTableScanNode->dataRequired == FUNC_DATA_REQUIRED_NOT_LOAD) &&
236,455,276✔
3142
                       (pTableScanNode->scan.node.pConditions == NULL) && (pTableScanNode->interval == 0);
118,172,403✔
3143

3144
  int32_t j = 0;
118,203,749✔
3145
  for (int32_t i = 0; i < pCond->numOfCols; ++i) {
729,851,621✔
3146
    STargetNode* pNode = (STargetNode*)nodesListGetNode(pTableScanNode->scan.pScanCols, i);
611,640,302✔
3147
    if (!pNode) {
611,511,562✔
3148
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
3149
      return terrno;
×
3150
    }
3151
    SColumnNode* pColNode = (SColumnNode*)pNode->pExpr;
611,511,562✔
3152
    if (pColNode->colType == COLUMN_TYPE_TAG) {
611,538,398✔
3153
      continue;
×
3154
    }
3155

3156
    pCond->colList[j].type = pColNode->node.resType.type;
611,597,023✔
3157
    pCond->colList[j].bytes = pColNode->node.resType.bytes;
611,592,570✔
3158
    pCond->colList[j].colId = pColNode->colId;
611,565,138✔
3159
    pCond->colList[j].pk = pColNode->isPk;
611,654,406✔
3160

3161
    pCond->pSlotList[j] = pNode->slotId;
611,644,705✔
3162
    j += 1;
611,647,872✔
3163
  }
3164

3165
  pCond->numOfCols = j;
118,270,961✔
3166

3167
  if (applyExtWin) {
118,266,740✔
3168
    if (NULL != pTableScanNode->pExtScanRange) {
97,216,290✔
3169
      pCond->type = TIMEWINDOW_RANGE_EXTERNAL;
1,625,418✔
3170
      code = getQueryExtWindow(&pCond->twindows, pTableScanNode->pExtScanRange, &pCond->twindows, pCond->extTwindows);
1,625,418✔
3171
    } else if (readHandle->extWinRangeValid) {
95,536,202✔
3172
      pCond->type = TIMEWINDOW_RANGE_EXTERNAL;
×
3173
      code = getQueryExtWindow(&pCond->twindows, &readHandle->extWinRange, &pCond->twindows, pCond->extTwindows);
×
3174
    }
3175
  }
3176
  
3177
  return code;
118,183,145✔
3178
}
3179

3180
int32_t initQueryTableDataCondWithColArray(SQueryTableDataCond* pCond, SQueryTableDataCond* pOrgCond,
6,612,245✔
3181
                                           const SReadHandle* readHandle, SArray* colArray) {
3182
  int32_t code = TSDB_CODE_SUCCESS;
6,612,245✔
3183
  int32_t lino = 0;
6,612,245✔
3184

3185
  pCond->order = TSDB_ORDER_ASC;
6,612,245✔
3186
  pCond->numOfCols = (int32_t)taosArrayGetSize(colArray);
6,612,245✔
3187

3188
  pCond->colList = taosMemoryCalloc(pCond->numOfCols, sizeof(SColumnInfo));
6,612,245✔
3189
  QUERY_CHECK_NULL(pCond->colList, code, lino, _return, terrno);
6,612,245✔
3190

3191
  pCond->pSlotList = taosMemoryMalloc(sizeof(int32_t) * pCond->numOfCols);
6,612,245✔
3192
  QUERY_CHECK_NULL(pCond->pSlotList, code, lino, _return, terrno);
6,612,245✔
3193

3194
  pCond->twindows = pOrgCond->twindows;
6,612,245✔
3195
  pCond->type = pOrgCond->type;
6,612,245✔
3196
  pCond->startVersion = -1;
6,612,245✔
3197
  pCond->endVersion = -1;
6,612,245✔
3198
  pCond->skipRollup = true;
6,612,245✔
3199
  pCond->notLoadData = false;
6,612,245✔
3200

3201
  for (int32_t i = 0; i < pCond->numOfCols; ++i) {
33,272,861✔
3202
    SColIdPair* pColPair = taosArrayGet(colArray, i);
26,660,616✔
3203
    QUERY_CHECK_NULL(pColPair, code, lino, _return, terrno);
26,660,616✔
3204

3205
    bool find = false;
26,660,616✔
3206
    for (int32_t j = 0; j < pOrgCond->numOfCols; ++j) {
167,602,329✔
3207
      if (pOrgCond->colList[j].colId == pColPair->vtbColId) {
167,602,329✔
3208
        pCond->colList[i].type = pOrgCond->colList[j].type;
26,660,616✔
3209
        pCond->colList[i].bytes = pOrgCond->colList[j].bytes;
26,660,616✔
3210
        pCond->colList[i].colId = pColPair->orgColId;
26,660,616✔
3211
        pCond->colList[i].pk = pOrgCond->colList[j].pk;
26,660,616✔
3212
        pCond->pSlotList[i] = i;
26,660,616✔
3213
        find = true;
26,660,616✔
3214
        break;
26,660,616✔
3215
      }
3216
    }
3217
    QUERY_CHECK_CONDITION(find, code, lino, _return, TSDB_CODE_NOT_FOUND);
26,660,616✔
3218
  }
3219

3220
  return code;
6,612,245✔
3221
_return:
×
3222
  qError("%s failed at line %d since %s", __func__, lino, tstrerror(terrno));
×
3223
  taosMemoryFreeClear(pCond->colList);
×
3224
  taosMemoryFreeClear(pCond->pSlotList);
×
3225
  return code;
×
3226
}
3227

3228
void cleanupQueryTableDataCond(SQueryTableDataCond* pCond) {
265,820,402✔
3229
  taosMemoryFreeClear(pCond->colList);
265,820,402✔
3230
  taosMemoryFreeClear(pCond->pSlotList);
265,829,026✔
3231
}
265,808,772✔
3232

3233
int32_t convertFillType(int32_t mode) {
2,058,942✔
3234
  int32_t type = TSDB_FILL_NONE;
2,058,942✔
3235
  switch (mode) {
2,058,942✔
3236
    case FILL_MODE_PREV:
107,688✔
3237
      type = TSDB_FILL_PREV;
107,688✔
3238
      break;
107,688✔
3239
    case FILL_MODE_NONE:
×
3240
      type = TSDB_FILL_NONE;
×
3241
      break;
×
3242
    case FILL_MODE_NULL:
130,941✔
3243
      type = TSDB_FILL_NULL;
130,941✔
3244
      break;
130,941✔
3245
    case FILL_MODE_NULL_F:
21,215✔
3246
      type = TSDB_FILL_NULL_F;
21,215✔
3247
      break;
21,215✔
3248
    case FILL_MODE_NEXT:
95,178✔
3249
      type = TSDB_FILL_NEXT;
95,178✔
3250
      break;
95,178✔
3251
    case FILL_MODE_VALUE:
148,427✔
3252
      type = TSDB_FILL_SET_VALUE;
148,427✔
3253
      break;
148,427✔
3254
    case FILL_MODE_VALUE_F:
4,278✔
3255
      type = TSDB_FILL_SET_VALUE_F;
4,278✔
3256
      break;
4,278✔
3257
    case FILL_MODE_LINEAR:
151,046✔
3258
      type = TSDB_FILL_LINEAR;
151,046✔
3259
      break;
151,046✔
3260
    case FILL_MODE_NEAR:
1,400,689✔
3261
      type = TSDB_FILL_NEAR;
1,400,689✔
3262
      break;
1,400,689✔
3263
    default:
×
3264
      type = TSDB_FILL_NONE;
×
3265
  }
3266

3267
  return type;
2,058,942✔
3268
}
3269

3270
void getInitialStartTimeWindow(SInterval* pInterval, TSKEY ts, STimeWindow* w, bool ascQuery) {
1,820,431,360✔
3271
  if (ascQuery) {
1,820,431,360✔
3272
    *w = getAlignQueryTimeWindow(pInterval, ts);
1,820,836,833✔
3273
  } else {
3274
    // the start position of the first time window in the endpoint that spreads beyond the queried last timestamp
3275
    *w = getAlignQueryTimeWindow(pInterval, ts);
61,271✔
3276

3277
    int64_t key = w->skey;
145,024✔
3278
    while (key < ts) {  // moving towards end
159,526✔
3279
      key = getNextTimeWindowStart(pInterval, key, TSDB_ORDER_ASC);
75,412✔
3280
      if (key > ts) {
75,412✔
3281
        break;
60,910✔
3282
      }
3283

3284
      w->skey = key;
14,502✔
3285
    }
3286
    w->ekey = taosTimeAdd(w->skey, pInterval->interval, pInterval->intervalUnit, pInterval->precision, NULL) - 1;
145,024✔
3287
  }
3288
}
1,821,444,276✔
3289

3290
static STimeWindow doCalculateTimeWindow(int64_t ts, SInterval* pInterval) {
26,254,960✔
3291
  STimeWindow w = {0};
26,254,960✔
3292

3293
  w.skey = taosTimeTruncate(ts, pInterval);
26,254,960✔
3294
  w.ekey = taosTimeGetIntervalEnd(w.skey, pInterval);
26,253,280✔
3295
  return w;
26,255,293✔
3296
}
3297

3298
STimeWindow getFirstQualifiedTimeWindow(int64_t ts, STimeWindow* pWindow, SInterval* pInterval, int32_t order) {
1,623,095✔
3299
  STimeWindow win = *pWindow;
1,623,095✔
3300
  STimeWindow save = win;
1,623,095✔
3301
  while (win.skey <= ts && win.ekey >= ts) {
9,240,652✔
3302
    save = win;
7,617,557✔
3303
    // get previous time window
3304
    getNextTimeWindow(pInterval, &win, order == TSDB_ORDER_DESC ? TSDB_ORDER_ASC : TSDB_ORDER_DESC);
7,617,557✔
3305
  }
3306

3307
  return save;
1,623,095✔
3308
}
3309

3310
// get the correct time window according to the handled timestamp
3311
// todo refactor
3312
STimeWindow getActiveTimeWindow(SDiskbasedBuf* pBuf, SResultRowInfo* pResultRowInfo, int64_t ts, SInterval* pInterval,
40,767,295✔
3313
                                int32_t order) {
3314
  STimeWindow w = {0};
40,767,295✔
3315
  if (pResultRowInfo->cur.pageId == -1) {  // the first window, from the previous stored value
40,768,688✔
3316
    getInitialStartTimeWindow(pInterval, ts, &w, (order != TSDB_ORDER_DESC));
1,469,132✔
3317
    return w;
1,469,369✔
3318
  }
3319

3320
  SResultRow* pRow = getResultRowByPos(pBuf, &pResultRowInfo->cur, false);
39,298,669✔
3321
  if (pRow) {
39,298,024✔
3322
    TAOS_SET_OBJ_ALIGNED(&w, pRow->win);
39,297,785✔
3323
  }
3324

3325
  // in case of typical time window, we can calculate time window directly.
3326
  if (w.skey > ts || w.ekey < ts) {
39,298,503✔
3327
    w = doCalculateTimeWindow(ts, pInterval);
26,255,097✔
3328
  }
3329

3330
  if (pInterval->interval != pInterval->sliding) {
39,298,699✔
3331
    // it is an sliding window query, in which sliding value is not equalled to
3332
    // interval value, and we need to find the first qualified time window.
3333
    w = getFirstQualifiedTimeWindow(ts, &w, pInterval, order);
1,623,095✔
3334
  }
3335

3336
  return w;
39,296,696✔
3337
}
3338

3339
TSKEY getNextTimeWindowStart(const SInterval* pInterval, TSKEY start, int32_t order) {
2,147,483,647✔
3340
  int32_t factor = GET_FORWARD_DIRECTION_FACTOR(order);
2,147,483,647✔
3341
  TSKEY   nextStart = taosTimeAdd(start, -1 * pInterval->offset, pInterval->offsetUnit, pInterval->precision, NULL);
2,147,483,647✔
3342
  nextStart = taosTimeAdd(nextStart, factor * pInterval->sliding, pInterval->slidingUnit, pInterval->precision, NULL);
2,147,483,647✔
3343
  nextStart = taosTimeAdd(nextStart, pInterval->offset, pInterval->offsetUnit, pInterval->precision, NULL);
2,147,483,647✔
3344
  return nextStart;
2,147,483,647✔
3345
}
3346

3347
void getNextTimeWindow(const SInterval* pInterval, STimeWindow* tw, int32_t order) {
2,147,483,647✔
3348
  tw->skey = getNextTimeWindowStart(pInterval, tw->skey, order);
2,147,483,647✔
3349
  tw->ekey = taosTimeAdd(tw->skey, pInterval->interval, pInterval->intervalUnit, pInterval->precision, NULL) - 1;
2,147,483,647✔
3350
}
2,147,483,647✔
3351

3352
bool hasLimitOffsetInfo(SLimitInfo* pLimitInfo) {
301,841,533✔
3353
  return (pLimitInfo->limit.limit != -1 || pLimitInfo->limit.offset != -1 || pLimitInfo->slimit.limit != -1 ||
601,516,309✔
3354
          pLimitInfo->slimit.offset != -1);
299,674,893✔
3355
}
3356

3357
bool hasSlimitOffsetInfo(SLimitInfo* pLimitInfo) {
×
3358
  return (pLimitInfo->slimit.limit != -1 || pLimitInfo->slimit.offset != -1);
×
3359
}
3360

3361
void initLimitInfo(const SNode* pLimit, const SNode* pSLimit, SLimitInfo* pLimitInfo) {
281,065,608✔
3362
  SLimit limit = {.limit = getLimit(pLimit), .offset = getOffset(pLimit)};
281,065,608✔
3363
  SLimit slimit = {.limit = getLimit(pSLimit), .offset = getOffset(pSLimit)};
281,037,238✔
3364

3365
  pLimitInfo->limit = limit;
281,035,375✔
3366
  pLimitInfo->slimit = slimit;
281,048,596✔
3367
  pLimitInfo->remainOffset = limit.offset;
281,060,824✔
3368
  pLimitInfo->remainGroupOffset = slimit.offset;
281,049,959✔
3369
  pLimitInfo->numOfOutputRows = 0;
281,040,790✔
3370
  pLimitInfo->numOfOutputGroups = 0;
281,035,103✔
3371
  pLimitInfo->currentGroupId = 0;
281,045,853✔
3372
}
281,038,426✔
3373

3374
void resetLimitInfoForNextGroup(SLimitInfo* pLimitInfo) {
46,119,555✔
3375
  pLimitInfo->numOfOutputRows = 0;
46,119,555✔
3376
  pLimitInfo->remainOffset = pLimitInfo->limit.offset;
46,126,939✔
3377
}
46,127,201✔
3378

3379
int32_t tableListGetSize(const STableListInfo* pTableList, int32_t* pRes) {
288,440,804✔
3380
  if (taosArrayGetSize(pTableList->pTableList) != taosHashGetSize(pTableList->map)) {
288,440,804✔
3381
    qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR));
×
3382
    return TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR;
×
3383
  }
3384
  (*pRes) = taosArrayGetSize(pTableList->pTableList);
288,473,190✔
3385
  return TSDB_CODE_SUCCESS;
288,439,925✔
3386
}
3387

3388
uint64_t tableListGetSuid(const STableListInfo* pTableList) { return pTableList->idInfo.suid; }
3,331,216✔
3389

3390
STableKeyInfo* tableListGetInfo(const STableListInfo* pTableList, int32_t index) {
147,918,557✔
3391
  if (taosArrayGetSize(pTableList->pTableList) == 0) {
147,918,557✔
3392
    return NULL;
3,645✔
3393
  }
3394

3395
  return taosArrayGet(pTableList->pTableList, index);
147,911,618✔
3396
}
3397

3398
int32_t tableListFind(const STableListInfo* pTableList, uint64_t uid, int32_t startIndex) {
9,750✔
3399
  int32_t numOfTables = taosArrayGetSize(pTableList->pTableList);
9,750✔
3400
  if (startIndex >= numOfTables) {
9,750✔
3401
    return -1;
×
3402
  }
3403

3404
  for (int32_t i = startIndex; i < numOfTables; ++i) {
91,307✔
3405
    STableKeyInfo* p = taosArrayGet(pTableList->pTableList, i);
91,307✔
3406
    if (!p) {
91,307✔
3407
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
3408
      return -1;
×
3409
    }
3410
    if (p->uid == uid) {
91,307✔
3411
      return i;
9,750✔
3412
    }
3413
  }
3414
  return -1;
×
3415
}
3416

3417
void tableListGetSourceTableInfo(const STableListInfo* pTableList, uint64_t* psuid, uint64_t* uid, int32_t* type) {
62,996✔
3418
  *psuid = pTableList->idInfo.suid;
62,996✔
3419
  *uid = pTableList->idInfo.uid;
62,996✔
3420
  *type = pTableList->idInfo.tableType;
62,996✔
3421
}
62,996✔
3422

3423
uint64_t tableListGetTableGroupId(const STableListInfo* pTableList, uint64_t tableUid) {
516,283,397✔
3424
  int32_t* slot = taosHashGet(pTableList->map, &tableUid, sizeof(tableUid));
516,283,397✔
3425
  if (slot == NULL) {
516,769,213✔
3426
    qDebug("table:%" PRIu64 " not found in table list", tableUid);
×
3427
    return -1;
×
3428
  }
3429

3430
  STableKeyInfo* pKeyInfo = taosArrayGet(pTableList->pTableList, *slot);
516,769,213✔
3431
  if (pKeyInfo == NULL) {
516,810,921✔
3432
    qDebug("table:%" PRIu64 " not found in table list", tableUid);
×
3433
    return -1;
×
3434
  }
3435
  return pKeyInfo->groupId;
516,810,921✔
3436
}
3437

3438
// TODO handle the group offset info, fix it, the rule of group output will be broken by this function
3439
// int32_t tableListRemoveTableInfo(STableListInfo* pTableList, uint64_t uid) {
3440
//   int32_t code = TSDB_CODE_SUCCESS;
3441
//   int32_t lino = 0;
3442

3443
//   int32_t* slot = taosHashGet(pTableList->map, &uid, sizeof(uid));
3444
//   if (slot == NULL) {
3445
//     qDebug("table:%" PRIu64 " not found in table list", uid);
3446
//     return 0;
3447
//   }
3448

3449
//   taosArrayRemove(pTableList->pTableList, *slot);
3450
//   code = taosHashRemove(pTableList->map, &uid, sizeof(uid));
3451

3452
//   _end:
3453
//   if (code != TSDB_CODE_SUCCESS) {
3454
//     qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
3455
//   } else {
3456
//     qDebug("uid:%" PRIu64 ", remove from table list", uid);
3457
//   }
3458

3459
//   return code;
3460
// }
3461

3462
int32_t tableListAddTableInfo(STableListInfo* pTableList, uint64_t uid, uint64_t gid) {
193,772✔
3463
  int32_t code = TSDB_CODE_SUCCESS;
193,772✔
3464
  int32_t lino = 0;
193,772✔
3465
  if (pTableList->map == NULL) {
193,772✔
3466
    pTableList->map = taosHashInit(32, taosGetDefaultHashFunction(TSDB_DATA_TYPE_BINARY), false, HASH_ENTRY_LOCK);
×
3467
    QUERY_CHECK_NULL(pTableList->map, code, lino, _end, terrno);
×
3468
  }
3469

3470
  STableKeyInfo keyInfo = {.uid = uid, .groupId = gid};
193,950✔
3471
  void*         p = taosHashGet(pTableList->map, &uid, sizeof(uid));
193,894✔
3472
  if (p != NULL) {
193,828✔
3473
    qInfo("table:%" PRId64 " already in tableIdList, ignore it", uid);
172✔
3474
    goto _end;
172✔
3475
  }
3476

3477
  void* tmp = taosArrayPush(pTableList->pTableList, &keyInfo);
193,656✔
3478
  QUERY_CHECK_NULL(tmp, code, lino, _end, terrno);
193,600✔
3479

3480
  int32_t slot = (int32_t)taosArrayGetSize(pTableList->pTableList) - 1;
193,600✔
3481
  code = taosHashPut(pTableList->map, &uid, sizeof(uid), &slot, sizeof(slot));
193,712✔
3482
  if (code != TSDB_CODE_SUCCESS) {
193,890✔
3483
    // we have checked the existence of uid in hash map above
3484
    QUERY_CHECK_CONDITION((code != TSDB_CODE_DUP_KEY), code, lino, _end, TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR);
×
3485
    taosArrayPopTailBatch(pTableList->pTableList, 1);  // let's pop the last element in the array list
×
3486
  }
3487

3488
_end:
194,062✔
3489
  if (code != TSDB_CODE_SUCCESS) {
193,884✔
3490
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
3491
  } else {
3492
    qDebug("uid:%" PRIu64 ", groupId:%" PRIu64 " added into table list, slot:%d, total:%d", uid, gid, slot, slot + 1);
193,884✔
3493
  }
3494

3495
  return code;
194,118✔
3496
}
3497

3498
int32_t tableListGetGroupList(const STableListInfo* pTableList, int32_t ordinalGroupIndex, STableKeyInfo** pKeyInfo,
98,839,577✔
3499
                              int32_t* size) {
3500
  int32_t totalGroups = tableListGetOutputGroups(pTableList);
98,839,577✔
3501
  int32_t numOfTables = 0;
98,853,971✔
3502
  int32_t code = tableListGetSize(pTableList, &numOfTables);
98,855,771✔
3503
  if (code != TSDB_CODE_SUCCESS) {
98,837,895✔
3504
    qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(code));
×
3505
    return code;
×
3506
  }
3507

3508
  if (ordinalGroupIndex < 0 || ordinalGroupIndex >= totalGroups) {
98,837,895✔
3509
    return TSDB_CODE_INVALID_PARA;
×
3510
  }
3511

3512
  // here handle two special cases:
3513
  // 1. only one group exists, and 2. one table exists for each group.
3514
  if (totalGroups == 1) {
98,837,895✔
3515
    *size = numOfTables;
98,472,727✔
3516
    *pKeyInfo = (*size == 0) ? NULL : taosArrayGet(pTableList->pTableList, 0);
98,466,903✔
3517
    return TSDB_CODE_SUCCESS;
98,477,433✔
3518
  } else if (totalGroups == numOfTables) {
365,849✔
3519
    *size = 1;
325,464✔
3520
    *pKeyInfo = taosArrayGet(pTableList->pTableList, ordinalGroupIndex);
325,464✔
3521
    return TSDB_CODE_SUCCESS;
324,950✔
3522
  }
3523

3524
  int32_t offset = pTableList->groupOffset[ordinalGroupIndex];
40,473✔
3525
  if (ordinalGroupIndex < totalGroups - 1) {
55,184✔
3526
    *size = pTableList->groupOffset[ordinalGroupIndex + 1] - offset;
40,602✔
3527
  } else {
3528
    *size = numOfTables - offset;
14,582✔
3529
  }
3530

3531
  *pKeyInfo = taosArrayGet(pTableList->pTableList, offset);
55,184✔
3532
  return TSDB_CODE_SUCCESS;
55,184✔
3533
}
3534

3535
int32_t tableListGetOutputGroups(const STableListInfo* pTableList) { return pTableList->numOfOuputGroups; }
289,577,903✔
3536

3537
bool oneTableForEachGroup(const STableListInfo* pTableList) { return pTableList->oneTableForEachGroup; }
621,908✔
3538

3539
STableListInfo* tableListCreate() {
124,233,070✔
3540
  STableListInfo* pListInfo = taosMemoryCalloc(1, sizeof(STableListInfo));
124,233,070✔
3541
  if (pListInfo == NULL) {
124,226,031✔
3542
    return NULL;
×
3543
  }
3544

3545
  pListInfo->remainGroups = NULL;
124,226,031✔
3546
  pListInfo->pTableList = taosArrayInit(4, sizeof(STableKeyInfo));
124,227,000✔
3547
  if (pListInfo->pTableList == NULL) {
124,248,413✔
3548
    goto _error;
×
3549
  }
3550

3551
  pListInfo->map = taosHashInit(1024, taosGetDefaultHashFunction(TSDB_DATA_TYPE_BIGINT), false, HASH_ENTRY_LOCK);
124,232,512✔
3552
  if (pListInfo->map == NULL) {
124,267,291✔
3553
    goto _error;
×
3554
  }
3555

3556
  pListInfo->numOfOuputGroups = 1;
124,265,070✔
3557
  return pListInfo;
124,264,276✔
3558

3559
_error:
×
3560
  tableListDestroy(pListInfo);
×
3561
  return NULL;
×
3562
}
3563

3564
void tableListDestroy(STableListInfo* pTableListInfo) {
133,276,486✔
3565
  if (pTableListInfo == NULL) {
133,276,486✔
3566
    return;
9,042,745✔
3567
  }
3568

3569
  taosArrayDestroy(pTableListInfo->pTableList);
124,233,741✔
3570
  taosMemoryFreeClear(pTableListInfo->groupOffset);
124,250,976✔
3571

3572
  taosHashCleanup(pTableListInfo->map);
124,247,467✔
3573
  taosHashCleanup(pTableListInfo->remainGroups);
124,256,944✔
3574
  pTableListInfo->pTableList = NULL;
124,264,723✔
3575
  pTableListInfo->map = NULL;
124,261,741✔
3576
  taosMemoryFree(pTableListInfo);
124,255,680✔
3577
}
3578

3579
void tableListClear(STableListInfo* pTableListInfo) {
145,272✔
3580
  if (pTableListInfo == NULL) {
145,272✔
3581
    return;
×
3582
  }
3583

3584
  taosArrayClear(pTableListInfo->pTableList);
145,272✔
3585
  taosHashClear(pTableListInfo->map);
145,285✔
3586
  taosHashClear(pTableListInfo->remainGroups);
145,440✔
3587
  taosMemoryFree(pTableListInfo->groupOffset);
145,440✔
3588
  pTableListInfo->numOfOuputGroups = 1;
145,440✔
3589
  pTableListInfo->oneTableForEachGroup = false;
145,440✔
3590
}
3591

3592
static int32_t orderbyGroupIdComparFn(const void* p1, const void* p2) {
496,725,251✔
3593
  STableKeyInfo* pInfo1 = (STableKeyInfo*)p1;
496,725,251✔
3594
  STableKeyInfo* pInfo2 = (STableKeyInfo*)p2;
496,725,251✔
3595

3596
  if (pInfo1->groupId == pInfo2->groupId) {
496,725,251✔
3597
    return 0;
480,573,837✔
3598
  } else {
3599
    return pInfo1->groupId < pInfo2->groupId ? -1 : 1;
16,153,062✔
3600
  }
3601
}
3602

3603
int32_t sortTableGroup(STableListInfo* pTableListInfo) {
17,128,026✔
3604
  int32_t code = TSDB_CODE_SUCCESS;
17,128,026✔
3605
  taosArraySort(pTableListInfo->pTableList, orderbyGroupIdComparFn);
17,128,026✔
3606
  int32_t size = taosArrayGetSize(pTableListInfo->pTableList);
17,134,660✔
3607
  if (size == 0) {
17,133,202✔
3608
    pTableListInfo->numOfOuputGroups = 0;
×
3609
    return code;
×
3610
  }
3611

3612
  SArray* pList = taosArrayInit(4, sizeof(int32_t));
17,133,202✔
3613
  if (!pList) {
17,134,515✔
3614
    code = terrno;
×
3615
    goto end;
×
3616
  }
3617

3618
  STableKeyInfo* pInfo = taosArrayGet(pTableListInfo->pTableList, 0);
17,134,515✔
3619
  if (pInfo == NULL) {
17,130,987✔
3620
    qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
3621
    code = terrno;
×
3622
    goto end;
×
3623
  }
3624
  uint64_t gid = pInfo->groupId;
17,130,987✔
3625

3626
  int32_t start = 0;
17,130,882✔
3627
  void*   tmp = taosArrayPush(pList, &start);
17,135,407✔
3628
  if (tmp == NULL) {
17,135,407✔
3629
    qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
3630
    code = terrno;
×
3631
    goto end;
×
3632
  }
3633

3634
  for (int32_t i = 1; i < size; ++i) {
116,909,143✔
3635
    pInfo = taosArrayGet(pTableListInfo->pTableList, i);
99,773,803✔
3636
    if (pInfo == NULL) {
99,773,477✔
3637
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
3638
      code = terrno;
×
3639
      goto end;
×
3640
    }
3641
    if (pInfo->groupId != gid) {
99,773,477✔
3642
      tmp = taosArrayPush(pList, &i);
3,517,089✔
3643
      if (tmp == NULL) {
3,517,089✔
3644
        qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
3645
        code = terrno;
×
3646
        goto end;
×
3647
      }
3648
      gid = pInfo->groupId;
3,517,089✔
3649
    }
3650
  }
3651

3652
  pTableListInfo->numOfOuputGroups = taosArrayGetSize(pList);
17,134,245✔
3653
  pTableListInfo->groupOffset = taosMemoryMalloc(sizeof(int32_t) * pTableListInfo->numOfOuputGroups);
17,134,568✔
3654
  if (pTableListInfo->groupOffset == NULL) {
17,133,082✔
3655
    code = terrno;
×
3656
    goto end;
×
3657
  }
3658

3659
  memcpy(pTableListInfo->groupOffset, taosArrayGet(pList, 0), sizeof(int32_t) * pTableListInfo->numOfOuputGroups);
17,133,528✔
3660

3661
end:
17,133,073✔
3662
  taosArrayDestroy(pList);
17,134,493✔
3663
  return code;
17,133,331✔
3664
}
3665

3666
int32_t buildGroupIdMapForAllTables(STableListInfo* pTableListInfo, SReadHandle* pHandle, SScanPhysiNode* pScanNode,
106,839,791✔
3667
                                    SNodeList* group, bool groupSort, uint8_t* digest, SStorageAPI* pAPI, SHashObj* groupIdMap) {
3668
  int32_t code = TSDB_CODE_SUCCESS;
106,839,791✔
3669

3670
  bool   groupByTbname = groupbyTbname(group);
106,839,791✔
3671
  size_t numOfTables = taosArrayGetSize(pTableListInfo->pTableList);
106,846,633✔
3672
  if (!numOfTables) {
106,847,393✔
3673
    return code;
6,180✔
3674
  }
3675
  qDebug("numOfTables:%zu, groupByTbname:%d, group:%p", numOfTables, groupByTbname, group);
106,841,213✔
3676
  if (group == NULL || groupByTbname) {
106,830,106✔
3677
    if (tsCountAlwaysReturnValue && QUERY_NODE_PHYSICAL_PLAN_TABLE_SCAN == nodeType(pScanNode) &&
104,965,480✔
3678
        ((STableScanPhysiNode*)pScanNode)->needCountEmptyTable) {
81,787,083✔
3679
      pTableListInfo->remainGroups =
4,250,876✔
3680
          taosHashInit(numOfTables, taosGetDefaultHashFunction(TSDB_DATA_TYPE_BIGINT), false, HASH_NO_LOCK);
4,250,876✔
3681
      if (pTableListInfo->remainGroups == NULL) {
4,250,876✔
3682
        return terrno;
×
3683
      }
3684

3685
      for (int i = 0; i < numOfTables; i++) {
19,874,853✔
3686
        STableKeyInfo* info = taosArrayGet(pTableListInfo->pTableList, i);
15,623,047✔
3687
        if (!info) {
15,624,132✔
3688
          qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
3689
          return terrno;
×
3690
        }
3691
        info->groupId = groupByTbname ? info->uid : 0;
15,624,132✔
3692
        int32_t tempRes = taosHashPut(pTableListInfo->remainGroups, &(info->groupId), sizeof(info->groupId),
15,624,442✔
3693
                                      &(info->uid), sizeof(info->uid));
15,624,287✔
3694
        if (tempRes != TSDB_CODE_SUCCESS && tempRes != TSDB_CODE_DUP_KEY) {
15,623,977✔
3695
          qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(tempRes));
×
3696
          return tempRes;
×
3697
        }
3698
      }
3699
    } else {
3700
      for (int32_t i = 0; i < numOfTables; i++) {
470,256,040✔
3701
        STableKeyInfo* info = taosArrayGet(pTableListInfo->pTableList, i);
369,549,729✔
3702
        if (!info) {
369,544,430✔
3703
          qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
3704
          return terrno;
×
3705
        }
3706
        info->groupId = groupByTbname ? info->uid : 0;
369,544,430✔
3707
        
3708
      }
3709
    }
3710
    if (groupIdMap && group != NULL){
104,958,117✔
3711
      getColInfoResultForGroupbyForStream(pHandle->vnode, group, pTableListInfo, pAPI, groupIdMap);
106,024✔
3712
    }
3713

3714
    pTableListInfo->oneTableForEachGroup = groupByTbname;
104,958,497✔
3715
    if (numOfTables == 1 && pTableListInfo->idInfo.tableType == TSDB_CHILD_TABLE) {
104,971,982✔
3716
      pTableListInfo->oneTableForEachGroup = true;
31,996,655✔
3717
    }
3718

3719
    if (groupSort && groupByTbname) {
104,972,210✔
3720
      taosArraySort(pTableListInfo->pTableList, orderbyGroupIdComparFn);
1,451,826✔
3721
      pTableListInfo->numOfOuputGroups = numOfTables;
1,451,826✔
3722
    } else if (groupByTbname && pScanNode->groupOrderScan) {
103,520,384✔
3723
      pTableListInfo->numOfOuputGroups = numOfTables;
31,358✔
3724
    } else {
3725
      pTableListInfo->numOfOuputGroups = 1;
103,489,026✔
3726
    }
3727
    if (groupSort || pScanNode->groupOrderScan) {
104,972,473✔
3728
      code = sortTableGroup(pTableListInfo);
16,986,417✔
3729
    }
3730
  } else {
3731
    bool initRemainGroups = false;
1,864,626✔
3732
    if (QUERY_NODE_PHYSICAL_PLAN_TABLE_SCAN == nodeType(pScanNode)) {
1,864,626✔
3733
      STableScanPhysiNode* pTableScanNode = (STableScanPhysiNode*)pScanNode;
1,682,437✔
3734
      if (tsCountAlwaysReturnValue && pTableScanNode->needCountEmptyTable &&
1,682,437✔
3735
          !(groupSort || pScanNode->groupOrderScan)) {
876,695✔
3736
        initRemainGroups = true;
850,235✔
3737
      }
3738
    }
3739

3740
    code = getColInfoResultForGroupby(pHandle->vnode, group, pTableListInfo, digest, pAPI, initRemainGroups, groupIdMap);
1,864,626✔
3741
    if (code != TSDB_CODE_SUCCESS) {
1,864,626✔
3742
      return code;
×
3743
    }
3744

3745
    if (pScanNode->groupOrderScan) pTableListInfo->numOfOuputGroups = taosArrayGetSize(pTableListInfo->pTableList);
1,864,626✔
3746

3747
    if (groupSort || pScanNode->groupOrderScan) {
1,864,626✔
3748
      code = sortTableGroup(pTableListInfo);
150,463✔
3749
    }
3750
  }
3751

3752
  // add all table entry in the hash map
3753
  size_t size = taosArrayGetSize(pTableListInfo->pTableList);
106,831,513✔
3754
  for (int32_t i = 0; i < size; ++i) {
501,113,552✔
3755
    STableKeyInfo* p = taosArrayGet(pTableListInfo->pTableList, i);
394,272,070✔
3756
    if (!p) {
394,251,300✔
3757
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
3758
      return terrno;
×
3759
    }
3760
    int32_t tempRes = taosHashPut(pTableListInfo->map, &p->uid, sizeof(uint64_t), &i, sizeof(int32_t));
394,251,300✔
3761
    if (tempRes != TSDB_CODE_SUCCESS && tempRes != TSDB_CODE_DUP_KEY) {
394,285,653✔
3762
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(tempRes));
×
3763
      return tempRes;
×
3764
    }
3765
  }
3766

3767
  return code;
106,841,193✔
3768
}
3769

3770
int32_t createScanTableListInfo(SScanPhysiNode* pScanNode, SNodeList* pGroupTags, bool groupSort, SReadHandle* pHandle,
118,648,224✔
3771
                                STableListInfo* pTableListInfo, SNode* pTagCond, SNode* pTagIndexCond,
3772
                                SExecTaskInfo* pTaskInfo, SHashObj* groupIdMap) {
3773
  int64_t     st = taosGetTimestampUs();
118,679,388✔
3774
  const char* idStr = GET_TASKID(pTaskInfo);
118,679,388✔
3775

3776
  if (pHandle == NULL) {
118,598,721✔
3777
    qError("invalid handle, in creating operator tree, %s", idStr);
×
3778
    return TSDB_CODE_INVALID_PARA;
×
3779
  }
3780

3781
  if (pHandle->uid != 0) {
118,598,721✔
3782
    pScanNode->uid = pHandle->uid;
32,872✔
3783
    pScanNode->tableType = TSDB_CHILD_TABLE;
32,872✔
3784
  }
3785
  uint8_t digest[17] = {0};
118,622,232✔
3786
  int32_t code = getTableList(pHandle->vnode, pScanNode, pTagCond, pTagIndexCond, pTableListInfo, digest, idStr,
118,631,877✔
3787
                              &pTaskInfo->storageAPI, pTaskInfo->pStreamRuntimeInfo);
118,648,039✔
3788
  if (code != TSDB_CODE_SUCCESS) {
118,723,983✔
3789
    qError("failed to getTableList, code:%s", tstrerror(code));
1,098✔
3790
    return code;
1,098✔
3791
  }
3792

3793
  int32_t numOfTables = taosArrayGetSize(pTableListInfo->pTableList);
118,722,885✔
3794

3795
  int64_t st1 = taosGetTimestampUs();
118,732,063✔
3796
  pTaskInfo->cost.extractListTime = (st1 - st) / 1000.0;
118,732,063✔
3797
  qDebug("extract queried table list completed, %d tables, elapsed time:%.2f ms %s", numOfTables,
118,729,249✔
3798
         pTaskInfo->cost.extractListTime, idStr);
3799

3800
  if (numOfTables == 0) {
118,728,816✔
3801
    qDebug("no table qualified for query, %s", idStr);
11,990,128✔
3802
    return TSDB_CODE_SUCCESS;
11,990,128✔
3803
  }
3804

3805
  code = buildGroupIdMapForAllTables(pTableListInfo, pHandle, pScanNode, pGroupTags, groupSort, digest, &pTaskInfo->storageAPI, groupIdMap);
106,738,688✔
3806
  if (code != TSDB_CODE_SUCCESS) {
106,741,410✔
3807
    return code;
×
3808
  }
3809

3810
  pTaskInfo->cost.groupIdMapTime = (taosGetTimestampUs() - st1) / 1000.0;
106,743,504✔
3811
  qDebug("generate group id map completed, elapsed time:%.2f ms %s", pTaskInfo->cost.groupIdMapTime, idStr);
106,740,016✔
3812

3813
  return TSDB_CODE_SUCCESS;
106,741,857✔
3814
}
3815

3816
char* getStreamOpName(uint16_t opType) {
9,292,249✔
3817
  switch (opType) {
9,292,249✔
3818
    case QUERY_NODE_PHYSICAL_PLAN_STREAM_SCAN:
×
3819
      return "stream scan";
×
3820
    case QUERY_NODE_PHYSICAL_PLAN_PROJECT:
9,055,919✔
3821
      return "project";
9,055,919✔
3822
    case QUERY_NODE_PHYSICAL_PLAN_EXTERNAL_WINDOW:
236,737✔
3823
      return "external window";
236,737✔
3824
  }
3825
  return "error name";
×
3826
}
3827

3828
void printDataBlock(SSDataBlock* pBlock, const char* flag, const char* taskIdStr, int64_t qId) {
248,070,096✔
3829
  if (qDebugFlag & DEBUG_TRACE) {
248,070,096✔
3830
    if (!pBlock) {
38,258✔
3831
      qDebug("%" PRIx64 " %s %s %s: Block is Null", qId, taskIdStr, flag, __func__);
6,919✔
3832
      return;
6,919✔
3833
    } else if (pBlock->info.rows == 0) {
31,339✔
3834
      qDebug("%" PRIx64 " %s %s %s: Block is Empty. block type %d", qId, taskIdStr, flag, __func__, pBlock->info.type);
×
3835
      return;
×
3836
    }
3837
    
3838
    char*   pBuf = NULL;
31,339✔
3839
    int32_t code = dumpBlockData(pBlock, flag, &pBuf, taskIdStr, qId);
31,339✔
3840
    if (code == 0) {
31,339✔
3841
      qDebugL("%" PRIx64 " %s %s", qId, __func__, pBuf);
31,339✔
3842
      taosMemoryFree(pBuf);
31,339✔
3843
    }
3844
  }
3845
}
3846

3847
void printSpecDataBlock(SSDataBlock* pBlock, const char* flag, const char* opStr, const char* taskIdStr) {
×
3848
  if (!pBlock) {
×
3849
    qDebug("%s===stream===%s %s: Block is Null", taskIdStr, flag, opStr);
×
3850
    return;
×
3851
  } else if (pBlock->info.rows == 0) {
×
3852
    qDebug("%s===stream===%s %s: Block is Empty. block type %d.skey:%" PRId64 ",ekey:%" PRId64 ",version%" PRId64,
×
3853
           taskIdStr, flag, opStr, pBlock->info.type, pBlock->info.window.skey, pBlock->info.window.ekey,
3854
           pBlock->info.version);
3855
    return;
×
3856
  }
3857
  if (qDebugFlag & DEBUG_TRACE) {
×
3858
    char* pBuf = NULL;
×
3859
    char  flagBuf[64];
×
3860
    snprintf(flagBuf, sizeof(flagBuf), "%s %s", flag, opStr);
×
3861
    int32_t code = dumpBlockData(pBlock, flagBuf, &pBuf, taskIdStr, 0);
×
3862
    if (code == 0) {
×
3863
      qDebug("%s", pBuf);
×
3864
      taosMemoryFree(pBuf);
×
3865
    }
3866
  }
3867
}
3868

3869
TSKEY getStartTsKey(STimeWindow* win, const TSKEY* tsCols) { return tsCols == NULL ? win->skey : tsCols[0]; }
8,920,672✔
3870

3871
void updateTimeWindowInfo(SColumnInfoData* pColData, const STimeWindow* pWin, int64_t delta) {
2,147,483,647✔
3872
  int64_t* ts = (int64_t*)pColData->pData;
2,147,483,647✔
3873

3874
  int64_t duration = pWin->ekey > pWin->skey ? pWin->ekey - pWin->skey + delta : pWin->skey - pWin->ekey + delta;
2,147,483,647✔
3875
  ts[2] = duration;            // set the duration
2,147,483,647✔
3876
  ts[3] = pWin->skey;          // window start key
2,147,483,647✔
3877
  ts[4] = pWin->ekey + delta;  // window end key
2,147,483,647✔
3878
}
2,147,483,647✔
3879

3880
int32_t compKeys(const SArray* pSortGroupCols, const char* oldkeyBuf, int32_t oldKeysLen, const SSDataBlock* pBlock,
933,019,468✔
3881
                 int32_t rowIndex) {
3882
  SColumnDataAgg* pColAgg = NULL;
933,019,468✔
3883
  const char*     isNull = oldkeyBuf;
933,019,468✔
3884
  const char*     p = oldkeyBuf + sizeof(int8_t) * pSortGroupCols->size;
933,019,468✔
3885

3886
  for (int32_t i = 0; i < pSortGroupCols->size; ++i) {
2,147,483,647✔
3887
    const SColumn*         pCol = (SColumn*)TARRAY_GET_ELEM(pSortGroupCols, i);
1,448,306,050✔
3888
    const SColumnInfoData* pColInfoData = TARRAY_GET_ELEM(pBlock->pDataBlock, pCol->slotId);
1,449,489,059✔
3889
    if (pBlock->pBlockAgg) pColAgg = &pBlock->pBlockAgg[pCol->slotId];
1,449,186,028✔
3890

3891
    if (colDataIsNull(pColInfoData, pBlock->info.rows, rowIndex, pColAgg)) {
2,147,483,647✔
3892
      if (isNull[i] != 1) return 1;
99,827,659✔
3893
    } else {
3894
      if (isNull[i] != 0) return 1;
1,349,997,384✔
3895
      const char* val = colDataGetData(pColInfoData, rowIndex);
1,349,155,722✔
3896
      if (pCol->type == TSDB_DATA_TYPE_JSON) {
1,349,311,100✔
3897
        int32_t len = getJsonValueLen(val);
×
3898
        if (memcmp(p, val, len) != 0) return 1;
×
3899
        p += len;
×
3900
      } else if (IS_VAR_DATA_TYPE(pCol->type)) {
1,349,083,583✔
3901
        if (IS_STR_DATA_BLOB(pCol->type)) {
457,262,187✔
3902
          if (memcmp(p, val, blobDataTLen(val)) != 0) return 1;
29,481✔
3903
          p += blobDataTLen(val);
×
3904
        } else {
3905
          if (memcmp(p, val, varDataTLen(val)) != 0) return 1;
457,156,902✔
3906
          p += varDataTLen(val);
450,109,001✔
3907
        }
3908
      } else {
3909
        if (0 != memcmp(p, val, pCol->bytes)) return 1;
892,179,158✔
3910
        p += pCol->bytes;
874,812,267✔
3911
      }
3912
    }
3913
  }
3914
  if ((int32_t)(p - oldkeyBuf) != oldKeysLen) return 1;
907,941,426✔
3915
  return 0;
907,928,624✔
3916
}
3917

3918
int32_t buildKeys(char* keyBuf, const SArray* pSortGroupCols, const SSDataBlock* pBlock, int32_t rowIndex) {
24,976,226✔
3919
  uint32_t        colNum = pSortGroupCols->size;
24,976,226✔
3920
  SColumnDataAgg* pColAgg = NULL;
25,011,586✔
3921
  char*           isNull = keyBuf;
25,011,586✔
3922
  char*           p = keyBuf + sizeof(int8_t) * colNum;
25,011,586✔
3923

3924
  for (int32_t i = 0; i < colNum; ++i) {
74,739,298✔
3925
    const SColumn*         pCol = (SColumn*)TARRAY_GET_ELEM(pSortGroupCols, i);
49,686,736✔
3926
    const SColumnInfoData* pColInfoData = TARRAY_GET_ELEM(pBlock->pDataBlock, pCol->slotId);
49,784,288✔
3927
    if (pCol->slotId > pBlock->pDataBlock->size) continue;
49,816,736✔
3928

3929
    if (pBlock->pBlockAgg) pColAgg = &pBlock->pBlockAgg[pCol->slotId];
49,797,184✔
3930

3931
    if (colDataIsNull(pColInfoData, pBlock->info.rows, rowIndex, pColAgg)) {
99,514,080✔
3932
      isNull[i] = 1;
1,792,688✔
3933
    } else {
3934
      isNull[i] = 0;
47,911,728✔
3935
      const char* val = colDataGetData(pColInfoData, rowIndex);
47,957,072✔
3936
      if (pCol->type == TSDB_DATA_TYPE_JSON) {
47,894,464✔
3937
        int32_t len = getJsonValueLen(val);
×
3938
        memcpy(p, val, len);
×
3939
        p += len;
×
3940
      } else if (IS_VAR_DATA_TYPE(pCol->type)) {
47,912,144✔
3941
        if (IS_STR_DATA_BLOB(pCol->type)) {
7,142,182✔
3942
          blobDataCopy(p, val);
10,608✔
3943
          p += blobDataTLen(val);
×
3944
        } else {
3945
          varDataCopy(p, val);
7,152,166✔
3946
          p += varDataTLen(val);
7,174,838✔
3947
        }
3948
      } else {
3949
        memcpy(p, val, pCol->bytes);
40,800,538✔
3950
        p += pCol->bytes;
40,783,482✔
3951
      }
3952
    }
3953
  }
3954
  return (int32_t)(p - keyBuf);
25,052,562✔
3955
}
3956

3957
uint64_t calcGroupId(char* pData, int32_t len) {
2,147,483,647✔
3958
  T_MD5_CTX context;
2,147,483,647✔
3959
  tMD5Init(&context);
2,147,483,647✔
3960
  tMD5Update(&context, (uint8_t*)pData, len);
2,147,483,647✔
3961
  tMD5Final(&context);
2,147,483,647✔
3962

3963
  // NOTE: only extract the initial 8 bytes of the final MD5 digest
3964
  uint64_t id = 0;
2,147,483,647✔
3965
  memcpy(&id, context.digest, sizeof(uint64_t));
2,147,483,647✔
3966
  if (0 == id) memcpy(&id, context.digest + 8, sizeof(uint64_t));
2,147,483,647✔
3967
  return id;
2,147,483,647✔
3968
}
3969

3970
SNodeList* makeColsNodeArrFromSortKeys(SNodeList* pSortKeys) {
41,392✔
3971
  SNode*     node;
3972
  SNodeList* ret = NULL;
41,392✔
3973
  FOREACH(node, pSortKeys) {
126,080✔
3974
    SOrderByExprNode* pSortKey = (SOrderByExprNode*)node;
84,688✔
3975
    int32_t           code = nodesListMakeAppend(&ret, pSortKey->pExpr);
84,688✔
3976
    if (code != TSDB_CODE_SUCCESS) {
84,688✔
3977
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(code));
×
3978
      terrno = code;
×
3979
      return NULL;
×
3980
    }
3981
  }
3982
  return ret;
41,392✔
3983
}
3984

3985
int32_t extractKeysLen(const SArray* keys, int32_t* pLen) {
41,392✔
3986
  int32_t code = TSDB_CODE_SUCCESS;
41,392✔
3987
  int32_t lino = 0;
41,392✔
3988
  int32_t len = 0;
41,392✔
3989
  int32_t keyNum = taosArrayGetSize(keys);
41,392✔
3990
  for (int32_t i = 0; i < keyNum; ++i) {
105,560✔
3991
    SColumn* pCol = (SColumn*)taosArrayGet(keys, i);
64,168✔
3992
    QUERY_CHECK_NULL(pCol, code, lino, _end, terrno);
64,168✔
3993
    len += pCol->bytes;
64,168✔
3994
  }
3995
  len += sizeof(int8_t) * keyNum;  // null flag
41,392✔
3996
  *pLen = len;
41,392✔
3997

3998
_end:
41,392✔
3999
  if (code != TSDB_CODE_SUCCESS) {
41,392✔
4000
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
4001
  }
4002
  return code;
41,392✔
4003
}
4004

4005
int32_t parseErrorMsgFromAnalyticServer(SJson* pJson, const char* pId) {
×
4006
  int32_t code = TSDB_CODE_ANA_ANODE_RETURN_ERROR;
×
4007
  if (pJson == NULL) {
×
4008
    return code;
×
4009
  }
4010

4011
  char    pMsg[1024] = {0};
×
4012
  int32_t ret = tjsonGetStringValue(pJson, "msg", pMsg);
×
4013

4014
  if (ret == 0) {
×
4015
    qError("%s failed to exec imputation, msg:%s", pId, pMsg);
×
4016
    if (strstr(pMsg, "white noise") != NULL) {
×
4017
      code = TSDB_CODE_ANA_WN_DATA;
×
4018
    } else if (strstr(pMsg, "white-noise") != NULL) {
×
4019
      code = TSDB_CODE_ANA_WN_DATA;
×
4020
    } else if (strstr(pMsg, "[Errno 111] Connection refused") != NULL) {
×
4021
      code = TSDB_CODE_ANA_ALGO_NOT_LOAD;
×
4022
    }
4023
  } else {
4024
    qError("%s failed to extract msg from server, unknown error", pId);
×
4025
  }
4026

4027
  return code;
×
4028
}
STATUS · Troubleshooting · Open an Issue · Sales · Support · CAREERS · ENTERPRISE · START FREE TRIAL · SCHEDULE DEMO
ANNOUNCEMENTS · TWITTER · TOS & SLA · Supported CI Services · What's a CI service? · Automated Testing

© 2026 Coveralls, Inc