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

taosdata / TDengine / #4895

23 Dec 2025 01:08PM UTC coverage: 65.513% (-0.2%) from 65.72%
#4895

push

travis-ci

web-flow
fix: mem leak (#34023)

6 of 9 new or added lines in 1 file covered. (66.67%)

7770 existing lines in 123 files now uncovered.

184705 of 281937 relevant lines covered (65.51%)

112009834.14 hits per line

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

75.65
/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
#include "trpc.h"
37

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

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

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

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

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

64
static int64_t getLimit(const SNode* pLimit) {
960,581,558✔
65
  return (NULL == pLimit || NULL == ((SLimitNode*)pLimit)->limit) ? -1 : ((SLimitNode*)pLimit)->limit->datum.i;
960,581,558✔
66
}
67
static int64_t getOffset(const SNode* pLimit) {
960,472,396✔
68
  return (NULL == pLimit || NULL == ((SLimitNode*)pLimit)->offset) ? -1 : ((SLimitNode*)pLimit)->offset->datum.i;
960,472,396✔
69
}
70
static void releaseColInfoData(void* pCol);
71

72
void initResultRowInfo(SResultRowInfo* pResultRowInfo) {
413,049,493✔
73
  pResultRowInfo->size = 0;
413,049,493✔
74
  pResultRowInfo->cur.pageId = -1;
413,087,210✔
75
}
413,146,000✔
76

77
void closeResultRow(SResultRow* pResultRow) { pResultRow->closed = true; }
6,108,747✔
78

79
void resetResultRow(SResultRow* pResultRow, size_t entrySize) {
424,225,863✔
80
  pResultRow->numOfRows = 0;
424,225,863✔
81
  pResultRow->closed = false;
424,228,608✔
82
  pResultRow->endInterp = false;
424,228,608✔
83
  pResultRow->startInterp = false;
424,228,608✔
84

85
  if (entrySize > 0) {
424,231,170✔
86
    memset(pResultRow->pEntryInfo, 0, entrySize);
424,231,353✔
87
  }
88
}
424,234,830✔
89

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

95
size_t getResultRowSize(SqlFunctionCtx* pCtx, int32_t numOfOutput) {
231,807,763✔
96
  int32_t rowSize = (numOfOutput * sizeof(SResultRowEntryInfo)) + sizeof(SResultRow);
231,807,763✔
97

98
  for (int32_t i = 0; i < numOfOutput; ++i) {
913,472,071✔
99
    rowSize += pCtx[i].resDataInfo.interBufSize;
681,760,558✔
100
  }
101

102
  return rowSize;
231,711,513✔
103
}
104

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

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

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

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

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

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

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

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

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

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

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

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

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

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

218
void cleanupGroupResInfo(SGroupResInfo* pGroupResInfo) {
106,378,444✔
219
  taosMemoryFreeClear(pGroupResInfo->pBuf);
106,378,444✔
220
  if (pGroupResInfo->freeItem) {
106,380,533✔
221
    //    taosArrayDestroy(pGroupResInfo->pRows);
222
    taosArrayDestroyEx(pGroupResInfo->pRows, freeEx);
×
223
    pGroupResInfo->freeItem = false;
×
224
    pGroupResInfo->pRows = NULL;
×
225
  } else {
226
    taosArrayDestroy(pGroupResInfo->pRows);
106,378,902✔
227
    pGroupResInfo->pRows = NULL;
106,376,806✔
228
  }
229
  pGroupResInfo->index = 0;
106,382,173✔
230
  pGroupResInfo->delIndex = 0;
106,382,776✔
231
}
106,381,576✔
232

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

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

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

251
static int32_t resultrowComparDesc(const void* p1, const void* p2) { return resultrowComparAsc(p2, p1); }
1,990,565,232✔
252

253
int32_t initGroupedResultInfo(SGroupResInfo* pGroupResInfo, SSHashObj* pHashmap, int32_t order) {
67,115,513✔
254
  int32_t code = TSDB_CODE_SUCCESS;
67,115,513✔
255
  int32_t lino = 0;
67,115,513✔
256
  if (pGroupResInfo->pRows != NULL) {
67,115,513✔
257
    taosArrayDestroy(pGroupResInfo->pRows);
3,324,879✔
258
  }
259
  if (pGroupResInfo->pBuf) {
67,117,690✔
260
    taosMemoryFree(pGroupResInfo->pBuf);
3,324,879✔
261
    pGroupResInfo->pBuf = NULL;
3,324,879✔
262
  }
263

264
  // extract the result rows information from the hash map
265
  int32_t size = tSimpleHashGetSize(pHashmap);
67,116,539✔
266

267
  void* pData = NULL;
67,114,697✔
268
  pGroupResInfo->pRows = taosArrayInit(size, POINTER_BYTES);
67,114,697✔
269
  QUERY_CHECK_NULL(pGroupResInfo->pRows, code, lino, _end, terrno);
67,111,078✔
270

271
  size_t  keyLen = 0;
67,108,340✔
272
  int32_t iter = 0;
67,114,351✔
273
  int64_t bufLen = 0, offset = 0;
67,115,754✔
274

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

281
  pGroupResInfo->pBuf = taosMemoryMalloc(bufLen);
67,091,329✔
282
  QUERY_CHECK_NULL(pGroupResInfo->pBuf, code, lino, _end, terrno);
67,116,115✔
283

284
  iter = 0;
67,115,319✔
285
  while ((pData = tSimpleHashIterate(pHashmap, pData, &iter)) != NULL) {
2,147,483,647✔
286
    void* key = tSimpleHashGetKey(pData, &keyLen);
2,147,483,647✔
287

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

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

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

299
  if (order == TSDB_ORDER_ASC || order == TSDB_ORDER_DESC) {
67,101,611✔
300
    __compar_fn_t fn = (order == TSDB_ORDER_ASC) ? resultrowComparAsc : resultrowComparDesc;
13,878,989✔
301
    size = POINTER_BYTES;
13,878,989✔
302
    taosSort(pGroupResInfo->pRows->pData, taosArrayGetSize(pGroupResInfo->pRows), size, fn);
13,878,989✔
303
  }
304

305
  pGroupResInfo->index = 0;
67,101,611✔
306

307
_end:
67,108,347✔
308
  if (code != TSDB_CODE_SUCCESS) {
67,114,903✔
309
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
310
  }
311
  return code;
67,114,903✔
312
}
313

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

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

325
bool hasRemainResults(SGroupResInfo* pGroupResInfo) {
234,952,519✔
326
  if (pGroupResInfo->pRows == NULL) {
234,952,519✔
327
    return false;
×
328
  }
329

330
  return pGroupResInfo->index < taosArrayGetSize(pGroupResInfo->pRows);
234,958,449✔
331
}
332

333
int32_t getNumOfTotalRes(SGroupResInfo* pGroupResInfo) {
126,427,497✔
334
  if (pGroupResInfo->pRows == 0) {
126,427,497✔
335
    return 0;
×
336
  }
337

338
  return (int32_t)taosArrayGetSize(pGroupResInfo->pRows);
126,433,198✔
339
}
340

341
SArray* createSortInfo(SNodeList* pNodeList) {
42,375,483✔
342
  size_t numOfCols = 0;
42,375,483✔
343

344
  if (pNodeList != NULL) {
42,375,483✔
345
    numOfCols = LIST_LENGTH(pNodeList);
42,324,523✔
346
  } else {
347
    numOfCols = 0;
51,108✔
348
  }
349

350
  SArray* pList = taosArrayInit(numOfCols, sizeof(SBlockOrderInfo));
42,374,005✔
351
  if (pList == NULL) {
42,369,219✔
352
    return pList;
×
353
  }
354

355
  for (int32_t i = 0; i < numOfCols; ++i) {
93,636,910✔
356
    SOrderByExprNode* pSortKey = (SOrderByExprNode*)nodesListGetNode(pNodeList, i);
51,260,831✔
357
    if (!pSortKey) {
51,265,672✔
358
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
359
      taosArrayDestroy(pList);
×
360
      pList = NULL;
×
361
      terrno = TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR;
×
362
      break;
×
363
    }
364
    SBlockOrderInfo bi = {0};
51,265,672✔
365
    bi.order = (pSortKey->order == ORDER_ASC) ? TSDB_ORDER_ASC : TSDB_ORDER_DESC;
51,262,445✔
366
    bi.nullFirst = (pSortKey->nullOrder == NULL_ORDER_FIRST);
51,264,519✔
367

368
    if (nodeType(pSortKey->pExpr) != QUERY_NODE_COLUMN) {
51,267,796✔
369
      qError("invalid order by expr type:%d", nodeType(pSortKey->pExpr));
×
370
      taosArrayDestroy(pList);
×
371
      pList = NULL;
×
372
      terrno = TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR;
×
373
      break;
×
374
    }
375
    
376
    SColumnNode* pColNode = (SColumnNode*)pSortKey->pExpr;
51,253,742✔
377
    bi.slotId = pColNode->slotId;
51,263,184✔
378
    void* tmp = taosArrayPush(pList, &bi);
51,267,691✔
379
    if (!tmp) {
51,267,691✔
380
      taosArrayDestroy(pList);
×
381
      pList = NULL;
×
382
      break;
×
383
    }
384
  }
385

386
  return pList;
42,375,383✔
387
}
388

389
SSDataBlock* createDataBlockFromDescNode(void* p) {
584,883,910✔
390
  SDataBlockDescNode* pNode = (SDataBlockDescNode*)p;
584,883,910✔
391
  int32_t      numOfCols = LIST_LENGTH(pNode->pSlots);
584,883,910✔
392
  SSDataBlock* pBlock = NULL;
584,840,159✔
393
  int32_t      code = createDataBlock(&pBlock);
584,800,331✔
394
  if (code) {
584,724,663✔
395
    terrno = code;
×
396
    return NULL;
×
397
  }
398

399
  pBlock->info.id.blockId = pNode->dataBlockId;
584,724,663✔
400
  pBlock->info.type = STREAM_INVALID;
584,696,565✔
401
  pBlock->info.calWin = (STimeWindow){.skey = INT64_MIN, .ekey = INT64_MAX};
584,837,814✔
402
  pBlock->info.watermark = INT64_MIN;
584,825,906✔
403

404
  for (int32_t i = 0; i < numOfCols; ++i) {
2,147,483,647✔
405
    SSlotDescNode* pDescNode = (SSlotDescNode*)nodesListGetNode(pNode->pSlots, i);
2,086,975,460✔
406
    if (!pDescNode) {
2,086,888,024✔
407
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(code));
×
408
      blockDataDestroy(pBlock);
×
409
      pBlock = NULL;
×
410
      terrno = TSDB_CODE_INVALID_PARA;
×
411
      break;
×
412
    }
413
    SColumnInfoData idata =
2,086,833,209✔
414
        createColumnInfoData(pDescNode->dataType.type, pDescNode->dataType.bytes, pDescNode->slotId);
2,087,052,468✔
415
    idata.info.scale = pDescNode->dataType.scale;
2,087,083,395✔
416
    idata.info.precision = pDescNode->dataType.precision;
2,087,115,321✔
417
    idata.info.noData = pDescNode->reserve;
2,087,146,951✔
418

419
    code = blockDataAppendColInfo(pBlock, &idata);
2,087,167,453✔
420
    if (code != TSDB_CODE_SUCCESS) {
2,087,123,174✔
421
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(code));
135,635✔
422
      blockDataDestroy(pBlock);
135,635✔
423
      pBlock = NULL;
×
424
      terrno = code;
×
425
      break;
×
426
    }
427
  }
428

429
  return pBlock;
585,059,063✔
430
}
431

432
int32_t prepareDataBlockBuf(SSDataBlock* pDataBlock, SColMatchInfo* pMatchInfo) {
211,833,144✔
433
  SDataBlockInfo* pBlockInfo = &pDataBlock->info;
211,833,144✔
434

435
  for (int32_t i = 0; i < taosArrayGetSize(pMatchInfo->pList); ++i) {
989,246,044✔
436
    SColMatchItem* pItem = taosArrayGet(pMatchInfo->pList, i);
784,279,077✔
437
    if (!pItem) {
784,334,342✔
438
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
439
      return terrno;
×
440
    }
441

442
    if (pItem->isPk) {
784,334,342✔
443
      SColumnInfoData* pInfoData = taosArrayGet(pDataBlock->pDataBlock, pItem->dstSlotId);
7,007,617✔
444
      if (!pInfoData) {
6,946,210✔
445
        qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
446
        return terrno;
×
447
      }
448
      pBlockInfo->pks[0].type = pInfoData->info.type;
6,946,210✔
449
      pBlockInfo->pks[1].type = pInfoData->info.type;
6,955,954✔
450

451
      // allocate enough buffer size, which is pInfoData->info.bytes
452
      if (IS_VAR_DATA_TYPE(pItem->dataType.type)) {
6,961,522✔
453
        pBlockInfo->pks[0].pData = taosMemoryCalloc(1, pInfoData->info.bytes);
2,323,990✔
454
        if (pBlockInfo->pks[0].pData == NULL) {
2,317,726✔
455
          return terrno;
×
456
        }
457

458
        pBlockInfo->pks[1].pData = taosMemoryCalloc(1, pInfoData->info.bytes);
2,320,510✔
459
        if (pBlockInfo->pks[1].pData == NULL) {
2,317,030✔
460
          taosMemoryFreeClear(pBlockInfo->pks[0].pData);
×
461
          return terrno;
×
462
        }
463

464
        pBlockInfo->pks[0].nData = pInfoData->info.bytes;
2,320,510✔
465
        pBlockInfo->pks[1].nData = pInfoData->info.bytes;
2,323,990✔
466
      }
467

468
      break;
6,955,258✔
469
    }
470
  }
471

472
  return TSDB_CODE_SUCCESS;
211,886,137✔
473
}
474

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

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

498
    res->translate = true;
123,756✔
499
    res->node.resType = pSColumnNode->node.resType;
123,756✔
500

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

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

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

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

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

559
  return DEAL_RES_CONTINUE;
433,146✔
560
}
561

562
int32_t isQualifiedTable(int64_t uid, SNode* pTagCond, void* vnode, bool* pQualified, SStorageAPI* pAPI) {
61,878✔
563
  int32_t     code = TSDB_CODE_SUCCESS;
61,878✔
564
  SMetaReader mr = {0};
61,878✔
565

566
  pAPI->metaReaderFn.initReader(&mr, vnode, META_READER_LOCK, &pAPI->metaFn);
61,878✔
567
  code = pAPI->metaReaderFn.getEntryGetUidCache(&mr, uid);
61,878✔
568
  if (TSDB_CODE_SUCCESS != code) {
61,878✔
569
    pAPI->metaReaderFn.clearReader(&mr);
×
570
    *pQualified = false;
×
571

572
    return TSDB_CODE_SUCCESS;
×
573
  }
574

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

592
  SNode* pNew = NULL;
61,878✔
593
  code = scalarCalculateConstants(pTagCondTmp, &pNew);
61,878✔
594
  if (TSDB_CODE_SUCCESS != code) {
61,878✔
595
    terrno = code;
×
596
    nodesDestroyNode(pTagCondTmp);
×
UNCOV
597
    *pQualified = false;
×
598

UNCOV
599
    return code;
×
600
  }
601

602
  SValueNode* pValue = (SValueNode*)pNew;
61,878✔
603
  *pQualified = pValue->datum.b;
61,878✔
604

605
  nodesDestroyNode(pNew);
61,878✔
606
  return TSDB_CODE_SUCCESS;
61,878✔
607
}
608

609
static EDealRes getColumn(SNode** pNode, void* pContext) {
46,951,893✔
610
  tagFilterAssist* pData = (tagFilterAssist*)pContext;
46,951,893✔
611
  SColumnNode*     pSColumnNode = NULL;
46,951,893✔
612
  if (QUERY_NODE_COLUMN == nodeType((*pNode))) {
46,953,369✔
613
    pSColumnNode = *(SColumnNode**)pNode;
15,452,469✔
614
  } else if (QUERY_NODE_FUNCTION == nodeType((*pNode))) {
31,508,982✔
615
    SFunctionNode* pFuncNode = *(SFunctionNode**)(pNode);
641,244✔
616
    if (pFuncNode->funcType == FUNCTION_TYPE_TBNAME) {
641,244✔
617
      pData->code = nodesMakeNode(QUERY_NODE_COLUMN, (SNode**)&pSColumnNode);
591,924✔
618
      if (NULL == pSColumnNode) {
591,924✔
UNCOV
619
        return DEAL_RES_ERROR;
×
620
      }
621
      pSColumnNode->colId = -1;
591,924✔
622
      pSColumnNode->colType = COLUMN_TYPE_TBNAME;
591,924✔
623
      pSColumnNode->node.resType.type = TSDB_DATA_TYPE_VARCHAR;
591,924✔
624
      pSColumnNode->node.resType.bytes = TSDB_TABLE_FNAME_LEN - 1 + VARSTR_HEADER_SIZE;
591,924✔
625
      nodesDestroyNode(*pNode);
591,924✔
626
      *pNode = (SNode*)pSColumnNode;
591,924✔
627
    } else {
628
      return DEAL_RES_CONTINUE;
49,320✔
629
    }
630
  } else {
631
    return DEAL_RES_CONTINUE;
30,864,733✔
632
  }
633

634
  void* data = taosHashGet(pData->colHash, &pSColumnNode->colId, sizeof(pSColumnNode->colId));
16,041,092✔
635
  if (!data) {
16,041,466✔
636
    int32_t tempRes =
637
        taosHashPut(pData->colHash, &pSColumnNode->colId, sizeof(pSColumnNode->colId), pNode, sizeof((*pNode)));
14,724,457✔
638
    if (tempRes != TSDB_CODE_SUCCESS && tempRes != TSDB_CODE_DUP_KEY) {
14,727,854✔
UNCOV
639
      return DEAL_RES_ERROR;
×
640
    }
641
    pSColumnNode->slotId = pData->index++;
14,727,854✔
642
    SColumnInfo cInfo = {.colId = pSColumnNode->colId,
14,727,044✔
643
                         .type = pSColumnNode->node.resType.type,
14,723,783✔
644
                         .bytes = pSColumnNode->node.resType.bytes,
14,724,444✔
645
                         .pk = pSColumnNode->isPk};
14,725,126✔
646
#if TAG_FILTER_DEBUG
647
    qDebug("tagfilter build column info, slotId:%d, colId:%d, type:%d", pSColumnNode->slotId, cInfo.colId, cInfo.type);
648
#endif
649
    void* tmp = taosArrayPush(pData->cInfoList, &cInfo);
14,725,883✔
650
    if (!tmp) {
14,725,794✔
UNCOV
651
      return DEAL_RES_ERROR;
×
652
    }
653
  } else {
654
    SColumnNode* col = *(SColumnNode**)data;
1,317,009✔
655
    pSColumnNode->slotId = col->slotId;
1,317,209✔
656
  }
657

658
  return DEAL_RES_CONTINUE;
16,042,180✔
659
}
660

661
static int32_t createResultData(SDataType* pType, int32_t numOfRows, SScalarParam* pParam) {
13,739,155✔
662
  SColumnInfoData* pColumnData = taosMemoryCalloc(1, sizeof(SColumnInfoData));
13,739,155✔
663
  if (pColumnData == NULL) {
13,737,559✔
UNCOV
664
    return terrno;
×
665
  }
666

667
  pColumnData->info.type = pType->type;
13,737,559✔
668
  pColumnData->info.bytes = pType->bytes;
13,737,698✔
669
  pColumnData->info.scale = pType->scale;
13,739,129✔
670
  pColumnData->info.precision = pType->precision;
13,737,981✔
671

672
  int32_t code = colInfoDataEnsureCapacity(pColumnData, numOfRows, true);
13,736,545✔
673
  if (code != TSDB_CODE_SUCCESS) {
13,738,926✔
674
    terrno = code;
×
675
    releaseColInfoData(pColumnData);
×
UNCOV
676
    return terrno;
×
677
  }
678

679
  pParam->columnData = pColumnData;
13,738,926✔
680
  pParam->colAlloced = true;
13,737,636✔
681
  return TSDB_CODE_SUCCESS;
13,734,269✔
682
}
683

684
static void releaseColInfoData(void* pCol) {
2,157,694✔
685
  if (pCol) {
2,157,694✔
686
    SColumnInfoData* col = (SColumnInfoData*)pCol;
2,157,694✔
687
    colDataDestroy(col);
2,157,694✔
688
    taosMemoryFree(col);
2,157,262✔
689
  }
690
}
2,156,717✔
691

692
void freeItem(void* p) {
190,844,201✔
693
  STUidTagInfo* pInfo = p;
190,844,201✔
694
  if (pInfo->pTagVal != NULL) {
190,844,201✔
695
    taosMemoryFree(pInfo->pTagVal);
190,426,567✔
696
  }
697
}
190,843,762✔
698

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

705
static int compareTagDataEntry(const void* a, const void* b) {
34,540✔
706
  STagDataEntry* p1 = (STagDataEntry*)a;
34,540✔
707
  STagDataEntry* p2 = (STagDataEntry*)b;
34,540✔
708
  return compareInt16Val(&p1->colId, &p2->colId);
34,540✔
709
}
710

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

726
    (void)memcpy(pStart, &entry->colId, sizeof(col_id_t));
34,540✔
727
    pStart += sizeof(col_id_t);
34,540✔
728

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

774
  return TSDB_CODE_SUCCESS;
17,270✔
775
}
776

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

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

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

802
  if (nodeType(pTagCond) == QUERY_NODE_OPERATOR) {
17,270✔
UNCOV
803
    extractTagDataEntry((SOperatorNode*)pTagCond, pIdWithVal);
×
804
  } else if (nodeType(pTagCond) == QUERY_NODE_LOGIC_CONDITION) {
17,270✔
805
    SNode* pChild = NULL;
17,270✔
806
    FOREACH(pChild, ((SLogicConditionNode*)pTagCond)->pParameterList) {
51,810✔
807
      extractTagDataEntry((SOperatorNode*)pChild, pIdWithVal);
34,540✔
808
    }
809
  }
810

811
  taosArraySort(pIdWithVal, compareTagDataEntry);
17,270✔
812

813
  return TSDB_CODE_SUCCESS;
17,270✔
814
}
815

816
static int32_t genStableTagFilterDigest(const SNode* pTagCond, T_MD5_CTX* pContext) {
17,270✔
817
  if (pTagCond == NULL) {
17,270✔
UNCOV
818
    return TSDB_CODE_SUCCESS;
×
819
  }
820

821
  char*   payload = NULL;
17,270✔
822
  int32_t len = 0;
17,270✔
823
  int32_t code = TSDB_CODE_SUCCESS;
17,270✔
824
  int32_t lino = 0;
17,270✔
825

826
  SArray* pIdWithVal = taosArrayInit(TARRAY_MIN_SIZE, sizeof(STagDataEntry));
17,270✔
827
  code = extractTagFilterTagDataEntries(pTagCond, pIdWithVal);
17,270✔
828
  QUERY_CHECK_CODE(code, lino, _end);
17,270✔
829
  for (int32_t i = 0; i < taosArrayGetSize(pIdWithVal); ++i) {
51,810✔
830
    STagDataEntry* pEntry = taosArrayGet(pIdWithVal, i);
34,540✔
831
    len += sizeof(col_id_t) + pEntry->bytes;
34,540✔
832
  }
833
  code = buildTagDataEntryKey(pIdWithVal, &payload, len);
17,270✔
834
  QUERY_CHECK_CODE(code, lino, _end);
17,270✔
835

836
  tMD5Init(pContext);
17,270✔
837
  tMD5Update(pContext, (uint8_t*)payload, (uint32_t)len);
17,270✔
838
  tMD5Final(pContext);
17,270✔
839

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

850
static int32_t genTagFilterDigest(const SNode* pTagCond, T_MD5_CTX* pContext) {
58,138✔
851
  if (pTagCond == NULL) {
58,138✔
852
    return TSDB_CODE_SUCCESS;
54,986✔
853
  }
854

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

863
  tMD5Init(pContext);
3,152✔
864
  tMD5Update(pContext, (uint8_t*)payload, (uint32_t)len);
3,152✔
865
  tMD5Final(pContext);
3,152✔
866

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

873
  taosMemoryFree(payload);
3,152✔
874
  return TSDB_CODE_SUCCESS;
3,152✔
875
}
876

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

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

892
  tMD5Init(pContext);
×
893
  tMD5Update(pContext, (uint8_t*)payload, (uint32_t)len);
×
UNCOV
894
  tMD5Final(pContext);
×
895

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

904
int32_t qGetColumnsFromNodeList(void* data, bool isList, SArray** pColList) {
13,555,269✔
905
  int32_t code = TSDB_CODE_SUCCESS;
13,555,269✔
906
  tagFilterAssist ctx = {0};
13,555,269✔
907
  ctx.colHash = taosHashInit(4, taosGetDefaultHashFunction(TSDB_DATA_TYPE_SMALLINT), false, HASH_NO_LOCK);
13,556,661✔
908
  if (ctx.colHash == NULL) {
13,556,349✔
909
    code = terrno;
×
UNCOV
910
    goto end;
×
911
  }
912

913
  ctx.index = 0;
13,556,349✔
914
  ctx.cInfoList = taosArrayInit(4, sizeof(SColumnInfo));
13,556,349✔
915
  if (ctx.cInfoList == NULL) {
13,556,991✔
916
    code = terrno;
1,496✔
UNCOV
917
    goto end;
×
918
  }
919

920
  if (isList) {
13,555,495✔
921
    SNode* pNode = NULL;
1,973,974✔
922
    FOREACH(pNode, (SNodeList*)data) {
4,138,052✔
923
      nodesRewriteExprPostOrder(&pNode, getColumn, (void*)&ctx);
2,164,603✔
924
      if (TSDB_CODE_SUCCESS != ctx.code) {
2,165,518✔
925
        code = ctx.code;
×
UNCOV
926
        goto end;
×
927
      }
928
      REPLACE_NODE(pNode);
2,165,518✔
929
    }
930
  } else {
931
    SNode* pNode = (SNode*)data;
11,581,521✔
932
    nodesRewriteExprPostOrder(&pNode, getColumn, (void*)&ctx);
11,582,342✔
933
    if (TSDB_CODE_SUCCESS != ctx.code) {
11,583,072✔
934
      code = ctx.code;
×
UNCOV
935
      goto end;
×
936
    }
937
  }
938
  
939
  if (pColList != NULL) *pColList = ctx.cInfoList;
13,557,643✔
940
  ctx.cInfoList = NULL;
13,556,258✔
941

942
end:
13,555,661✔
943
  taosHashCleanup(ctx.colHash);
13,556,053✔
944
  taosArrayDestroy(ctx.cInfoList);
13,552,235✔
945
  return code;
13,553,622✔
946
}
947

948
static int32_t buildGroupInfo(SColumnInfoData* pValue, int32_t i, SArray* gInfo) {
692,171✔
949
  int32_t code = TSDB_CODE_SUCCESS;
692,171✔
950
  SStreamGroupValue* v = taosArrayReserve(gInfo, 1);
692,171✔
951
  if (v == NULL) {
692,584✔
952
    code = terrno;
×
UNCOV
953
    goto end;
×
954
  }
955
  if (colDataIsNull_s(pValue, i)) {
1,384,854✔
956
    v->isNull = true;
12,560✔
957
  } else {
958
    v->isNull = false;
679,710✔
959
    char* data = colDataGetData(pValue, i);
680,024✔
960
    if (pValue->info.type == TSDB_DATA_TYPE_JSON) {
680,024✔
961
      if (tTagIsJson(data)) {
×
962
        code = TSDB_CODE_QRY_JSON_IN_GROUP_ERROR;
×
UNCOV
963
        goto end;
×
964
      }
965
      if (tTagIsJsonNull(data)) {
×
966
        v->isNull = true;
×
UNCOV
967
        goto end;
×
968
      }
969
      int32_t len = getJsonValueLen(data);
×
970
      v->data.type = pValue->info.type;
×
971
      v->data.nData = len;
×
972
      v->data.pData = taosMemoryCalloc(1, len + 1);
×
973
      if (v->data.pData == NULL) {
×
974
        code = terrno;
×
UNCOV
975
        goto end;
×
976
      }
977
      memcpy(v->data.pData, data, len);
×
UNCOV
978
      qDebug("buildGroupInfo:%d add json data len:%d, data:%s", i, len, (char*)v->data.pData);
×
979
    } else if (IS_VAR_DATA_TYPE(pValue->info.type)) {
679,611✔
980
      if (varDataTLen(data) > pValue->info.bytes) {
448,105✔
981
        code = TSDB_CODE_TDB_INVALID_TABLE_SCHEMA_VER;
×
UNCOV
982
        goto end;
×
983
      }
984
      v->data.type = pValue->info.type;
448,541✔
985
      v->data.nData = varDataLen(data);
448,541✔
986
      v->data.pData = taosMemoryCalloc(1, varDataLen(data) + 1);
448,954✔
987
      if (v->data.pData == NULL) {
448,541✔
988
        code = terrno;
×
UNCOV
989
        goto end;
×
990
      }
991
      memcpy(v->data.pData, varDataVal(data), varDataLen(data));
448,954✔
992
      qDebug("buildGroupInfo:%d add var data type:%d, len:%d, data:%s", i, pValue->info.type, varDataLen(data), (char*)v->data.pData);
448,954✔
993
    } else if (pValue->info.type == TSDB_DATA_TYPE_DECIMAL) {  // reader todo decimal
231,070✔
994
      v->data.type = pValue->info.type;
×
995
      v->data.nData = pValue->info.bytes;
×
996
      v->data.pData = taosMemoryCalloc(1, pValue->info.bytes);
×
997
      if (v->data.pData == NULL) {
×
998
        code = terrno;
×
UNCOV
999
        goto end;
×
1000
      }
1001
      memcpy(&v->data.pData, data, pValue->info.bytes);
×
UNCOV
1002
      qDebug("buildGroupInfo:%d add data type:%d, data:%"PRId64, i, pValue->info.type, v->data.val);
×
1003
    } else {  // reader todo decimal
1004
      v->data.type = pValue->info.type;
230,756✔
1005
      memcpy(&v->data.val, data, pValue->info.bytes);
231,070✔
1006
      qDebug("buildGroupInfo:%d add data type:%d, data:%"PRId64, i, pValue->info.type, v->data.val);
231,070✔
1007
    }
1008
  }
1009
end:
79,460✔
1010
  if (code != TSDB_CODE_SUCCESS) {
692,959✔
1011
    qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(code));
×
UNCOV
1012
    v->isNull = true;
×
1013
  }
1014
  return code;
692,959✔
1015
}
1016

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

1028
  int32_t rows = taosArrayGetSize(pTableListInfo->pTableList);
103,688✔
1029
  if (rows == 0) {
103,688✔
UNCOV
1030
    return;
×
1031
  }
1032

1033
  pUidTagList = taosArrayInit(8, sizeof(STUidTagInfo));
103,688✔
1034
  QUERY_CHECK_NULL(pUidTagList, code, lino, end, terrno);
103,688✔
1035

1036
  for (int32_t i = 0; i < rows; ++i) {
442,154✔
1037
    STableKeyInfo* pkeyInfo = taosArrayGet(pTableListInfo->pTableList, i);
338,466✔
1038
    QUERY_CHECK_NULL(pkeyInfo, code, lino, end, terrno);
338,466✔
1039
    STUidTagInfo info = {.uid = pkeyInfo->uid};
338,466✔
1040
    void*        tmp = taosArrayPush(pUidTagList, &info);
338,466✔
1041
    QUERY_CHECK_NULL(tmp, code, lino, end, terrno);
338,466✔
1042
  }
1043
 
1044
  if (taosArrayGetSize(pUidTagList) > 0) {
103,688✔
1045
    code = pAPI->metaFn.getTableTagsByUid(pVnode, pTableListInfo->idInfo.suid, pUidTagList);
103,688✔
1046
  } else {
UNCOV
1047
    code = pAPI->metaFn.getTableTags(pVnode, pTableListInfo->idInfo.suid, pUidTagList);
×
1048
  }
1049
  if (code != TSDB_CODE_SUCCESS) {
103,688✔
UNCOV
1050
    goto end;
×
1051
  }
1052

1053
  SArray* pColList = NULL;
103,688✔
1054
  code = qGetColumnsFromNodeList(group, true, &pColList);
103,688✔
1055
  if (code != TSDB_CODE_SUCCESS) {
103,688✔
UNCOV
1056
    goto end;
×
1057
  }
1058

1059
  for (int32_t i = 0; i < taosArrayGetSize(pColList); ++i) {
274,184✔
1060
    SColumnInfo* tmp = (SColumnInfo*)taosArrayGet(pColList, i);
170,496✔
1061
    if (tmp != NULL && tmp->colId == -1) {
170,496✔
1062
      tbNameIndex = i;
103,688✔
1063
    }
1064
  }
1065
  
1066
  int32_t numOfTables = taosArrayGetSize(pUidTagList);
103,688✔
1067
  pResBlock = createTagValBlockForFilter(pColList, numOfTables, pUidTagList, pVnode, pAPI);
103,688✔
1068
  taosArrayDestroy(pColList);
103,688✔
1069
  if (pResBlock == NULL) {
103,688✔
1070
    code = terrno;
×
UNCOV
1071
    goto end;
×
1072
  }
1073

1074
  pBlockList = taosArrayInit(2, POINTER_BYTES);
103,688✔
1075
  QUERY_CHECK_NULL(pBlockList, code, lino, end, terrno);
103,688✔
1076

1077
  void* tmp = taosArrayPush(pBlockList, &pResBlock);
103,688✔
1078
  QUERY_CHECK_NULL(tmp, code, lino, end, terrno);
103,688✔
1079

1080
  groupData = taosArrayInit(2, POINTER_BYTES);
103,688✔
1081
  QUERY_CHECK_NULL(groupData, code, lino, end, terrno);
103,688✔
1082

1083
  SNode* pNode = NULL;
103,688✔
1084
  FOREACH(pNode, group) {
274,184✔
1085
    SScalarParam output = {0};
170,496✔
1086

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

1101
      default:
×
1102
        code = TSDB_CODE_OPS_NOT_SUPPORT;
×
UNCOV
1103
        goto end;
×
1104
    }
1105

1106
    if (nodeType(pNode) == QUERY_NODE_COLUMN) {
170,496✔
1107
      SColumnNode*     pSColumnNode = (SColumnNode*)pNode;
170,496✔
1108
      SColumnInfoData* pColInfo = (SColumnInfoData*)taosArrayGet(pResBlock->pDataBlock, pSColumnNode->slotId);
170,496✔
1109
      QUERY_CHECK_NULL(pColInfo, code, lino, end, terrno);
170,496✔
1110
      code = colDataAssign(output.columnData, pColInfo, rows, NULL);
170,496✔
1111
    } else if (nodeType(pNode) == QUERY_NODE_VALUE) {
×
UNCOV
1112
      continue;
×
1113
    } else {
1114
      gTaskScalarExtra.pStreamInfo = NULL;
×
1115
      gTaskScalarExtra.pStreamRange = NULL;
×
UNCOV
1116
      code = scalarCalculate(pNode, pBlockList, &output, &gTaskScalarExtra);
×
1117
    }
1118

1119
    if (code != TSDB_CODE_SUCCESS) {
170,496✔
1120
      releaseColInfoData(output.columnData);
×
UNCOV
1121
      goto end;
×
1122
    }
1123

1124
    void* tmp = taosArrayPush(groupData, &output.columnData);
170,496✔
1125
    QUERY_CHECK_NULL(tmp, code, lino, end, terrno);
170,496✔
1126
  }
1127

1128
  for (int i = 0; i < rows; i++) {
442,154✔
1129
    gInfo = taosArrayInit(taosArrayGetSize(groupData), sizeof(SStreamGroupValue));
338,466✔
1130
    QUERY_CHECK_NULL(gInfo, code, lino, end, terrno);
338,466✔
1131

1132
    STableKeyInfo* info = taosArrayGet(pTableListInfo->pTableList, i);
338,466✔
1133
    QUERY_CHECK_NULL(info, code, lino, end, terrno);
338,466✔
1134

1135
    for (int j = 0; j < taosArrayGetSize(groupData); j++) {
850,665✔
1136
      SColumnInfoData* pValue = (SColumnInfoData*)taosArrayGetP(groupData, j);
512,199✔
1137
        int32_t ret = buildGroupInfo(pValue, i, gInfo);
512,199✔
1138
        if (ret != TSDB_CODE_SUCCESS) {
512,199✔
1139
          qError("buildGroupInfo failed at line %d since %s", __LINE__, tstrerror(ret));
×
UNCOV
1140
          goto end;
×
1141
        }
1142
        if (j == tbNameIndex) {
512,199✔
1143
          SStreamGroupValue* v = taosArrayGetLast(gInfo);
338,466✔
1144
          if (v != NULL){
338,466✔
1145
            v->isTbname = true;
338,466✔
1146
            v->uid = info->uid;
338,466✔
1147
          }
1148
        }
1149
    }
1150

1151
    int32_t ret = taosHashPut(groupIdMap, &info->uid, sizeof(info->uid), &gInfo, POINTER_BYTES);
338,466✔
1152
    if (ret != TSDB_CODE_SUCCESS) {
338,466✔
1153
      qError("put groupid to map failed at line %d since %s", __LINE__, tstrerror(ret));
×
UNCOV
1154
      goto end;
×
1155
    }
1156
    qDebug("put groupid to map gid:%" PRIu64, info->uid);
338,466✔
1157
    gInfo = NULL;
338,466✔
1158
  }
1159

1160
end:
103,688✔
1161
  blockDataDestroy(pResBlock);
103,688✔
1162
  taosArrayDestroy(pBlockList);
103,688✔
1163
  taosArrayDestroyEx(pUidTagList, freeItem);
103,688✔
1164
  taosArrayDestroyP(groupData, releaseColInfoData);
103,688✔
1165
  taosArrayDestroyEx(gInfo, tDestroySStreamGroupValue);
103,688✔
1166

1167
  if (code != TSDB_CODE_SUCCESS) {
103,688✔
UNCOV
1168
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
1169
  }
1170
}
1171

1172
int32_t getColInfoResultForGroupby(void* pVnode, SNodeList* group, STableListInfo* pTableListInfo, uint8_t* digest,
1,870,023✔
1173
                                   SStorageAPI* pAPI, bool initRemainGroups, SHashObj* groupIdMap) {
1174
  int32_t      code = TSDB_CODE_SUCCESS;
1,870,023✔
1175
  int32_t      lino = 0;
1,870,023✔
1176
  SArray*      pBlockList = NULL;
1,870,023✔
1177
  SSDataBlock* pResBlock = NULL;
1,870,023✔
1178
  void*        keyBuf = NULL;
1,870,436✔
1179
  SArray*      groupData = NULL;
1,870,436✔
1180
  SArray*      pUidTagList = NULL;
1,870,436✔
1181
  SArray*      tableList = NULL;
1,870,436✔
1182
  SArray*      gInfo = NULL;
1,870,436✔
1183

1184
  int32_t rows = taosArrayGetSize(pTableListInfo->pTableList);
1,870,436✔
1185
  if (rows == 0) {
1,870,436✔
UNCOV
1186
    return TSDB_CODE_SUCCESS;
×
1187
  } 
1188

1189
  T_MD5_CTX context = {0};
1,870,436✔
1190
  if (tsTagFilterCache && groupIdMap == NULL) {
1,870,436✔
1191
    SNodeListNode* listNode = NULL;
×
1192
    code = nodesMakeNode(QUERY_NODE_NODE_LIST, (SNode**)&listNode);
×
1193
    if (TSDB_CODE_SUCCESS != code) {
×
UNCOV
1194
      goto end;
×
1195
    }
1196
    listNode->pNodeList = group;
×
1197
    code = genTbGroupDigest((SNode*)listNode, digest, &context);
×
UNCOV
1198
    QUERY_CHECK_CODE(code, lino, end);
×
1199

UNCOV
1200
    nodesFree(listNode);
×
1201

UNCOV
1202
    code = pAPI->metaFn.metaGetCachedTbGroup(pVnode, pTableListInfo->idInfo.suid, context.digest,
×
1203
                                             tListLen(context.digest), &tableList);
UNCOV
1204
    QUERY_CHECK_CODE(code, lino, end);
×
1205

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

1215
  pUidTagList = taosArrayInit(8, sizeof(STUidTagInfo));
1,870,436✔
1216
  QUERY_CHECK_NULL(pUidTagList, code, lino, end, terrno);
1,870,436✔
1217

1218
  for (int32_t i = 0; i < rows; ++i) {
11,062,867✔
1219
    STableKeyInfo* pkeyInfo = taosArrayGet(pTableListInfo->pTableList, i);
9,192,431✔
1220
    QUERY_CHECK_NULL(pkeyInfo, code, lino, end, terrno);
9,192,431✔
1221
    STUidTagInfo info = {.uid = pkeyInfo->uid};
9,192,431✔
1222
    void*        tmp = taosArrayPush(pUidTagList, &info);
9,192,431✔
1223
    QUERY_CHECK_NULL(tmp, code, lino, end, terrno);
9,192,431✔
1224
  }
1225

1226
  if (taosArrayGetSize(pUidTagList) > 0) {
1,870,436✔
1227
    code = pAPI->metaFn.getTableTagsByUid(pVnode, pTableListInfo->idInfo.suid, pUidTagList);
1,870,436✔
1228
  } else {
UNCOV
1229
    code = pAPI->metaFn.getTableTags(pVnode, pTableListInfo->idInfo.suid, pUidTagList);
×
1230
  }
1231
  if (code != TSDB_CODE_SUCCESS) {
1,870,436✔
UNCOV
1232
    goto end;
×
1233
  }
1234

1235
  SArray* pColList = NULL;
1,870,436✔
1236
  code = qGetColumnsFromNodeList(group, true, &pColList); 
1,870,436✔
1237
  if (code != TSDB_CODE_SUCCESS) {
1,870,286✔
UNCOV
1238
    goto end;
×
1239
  }
1240

1241
  int32_t numOfTables = taosArrayGetSize(pUidTagList);
1,870,286✔
1242
  pResBlock = createTagValBlockForFilter(pColList, numOfTables, pUidTagList, pVnode, pAPI);
1,870,286✔
1243
  taosArrayDestroy(pColList);
1,870,436✔
1244
  if (pResBlock == NULL) {
1,870,436✔
1245
    code = terrno;
×
UNCOV
1246
    goto end;
×
1247
  }
1248

1249
  //  int64_t st1 = taosGetTimestampUs();
1250
  //  qDebug("generate tag block rows:%d, cost:%ld us", rows, st1-st);
1251

1252
  pBlockList = taosArrayInit(2, POINTER_BYTES);
1,870,436✔
1253
  QUERY_CHECK_NULL(pBlockList, code, lino, end, terrno);
1,870,436✔
1254

1255
  void* tmp = taosArrayPush(pBlockList, &pResBlock);
1,870,436✔
1256
  QUERY_CHECK_NULL(tmp, code, lino, end, terrno);
1,870,436✔
1257

1258
  groupData = taosArrayInit(2, POINTER_BYTES);
1,870,436✔
1259
  QUERY_CHECK_NULL(groupData, code, lino, end, terrno);
1,870,436✔
1260

1261
  SNode* pNode = NULL;
1,870,286✔
1262
  FOREACH(pNode, group) {
3,865,458✔
1263
    SScalarParam output = {0};
1,994,647✔
1264

1265
    switch (nodeType(pNode)) {
1,994,647✔
1266
      case QUERY_NODE_VALUE:
×
UNCOV
1267
        break;
×
1268
      case QUERY_NODE_COLUMN:
1,987,756✔
1269
      case QUERY_NODE_OPERATOR:
1270
      case QUERY_NODE_FUNCTION: {
1271
        SExprNode* expNode = (SExprNode*)pNode;
1,987,756✔
1272
        code = createResultData(&expNode->resType, rows, &output);
1,987,756✔
1273
        if (code != TSDB_CODE_SUCCESS) {
1,986,533✔
UNCOV
1274
          goto end;
×
1275
        }
1276
        break;
1,986,533✔
1277
      }
1278
      case QUERY_NODE_REMOTE_VALUE: {
7,416✔
1279
        SRemoteValueNode* pRemote = (SRemoteValueNode*)pNode;
7,416✔
1280
        code = qFetchRemoteValue(gTaskScalarExtra.pSubJobCtx, pRemote->subQIdx, pRemote);
7,416✔
1281
        QUERY_CHECK_CODE(code, lino, end);
7,416✔
1282
        break;
7,416✔
1283
      }
1284
      
1285
      default:
×
1286
        code = TSDB_CODE_OPS_NOT_SUPPORT;
×
UNCOV
1287
        goto end;
×
1288
    }
1289

1290
    if (nodeType(pNode) == QUERY_NODE_COLUMN) {
1,993,949✔
1291
      SColumnNode*     pSColumnNode = (SColumnNode*)pNode;
1,969,478✔
1292
      SColumnInfoData* pColInfo = (SColumnInfoData*)taosArrayGet(pResBlock->pDataBlock, pSColumnNode->slotId);
1,969,478✔
1293
      QUERY_CHECK_NULL(pColInfo, code, lino, end, terrno);
1,968,803✔
1294
      code = colDataAssign(output.columnData, pColInfo, rows, NULL);
1,968,803✔
1295
    } else if (nodeType(pNode) == QUERY_NODE_VALUE) {
24,321✔
1296
      continue;
7,416✔
1297
    } else {
1298
      gTaskScalarExtra.pStreamInfo = NULL;
17,587✔
1299
      gTaskScalarExtra.pStreamRange = NULL;
17,587✔
1300
      code = scalarCalculate(pNode, pBlockList, &output, &gTaskScalarExtra);
17,587✔
1301
    }
1302

1303
    if (code != TSDB_CODE_SUCCESS) {
1,987,231✔
1304
      releaseColInfoData(output.columnData);
×
UNCOV
1305
      goto end;
×
1306
    }
1307

1308
    void* tmp = taosArrayPush(groupData, &output.columnData);
1,987,756✔
1309
    QUERY_CHECK_NULL(tmp, code, lino, end, terrno);
1,987,756✔
1310
  }
1311

1312
  int32_t keyLen = 0;
1,869,596✔
1313
  SNode*  node;
1314
  FOREACH(node, group) {
3,862,010✔
1315
    SExprNode* pExpr = (SExprNode*)node;
1,993,957✔
1316
    keyLen += pExpr->resType.bytes;
1,993,957✔
1317
  }
1318

1319
  int32_t nullFlagSize = sizeof(int8_t) * LIST_LENGTH(group);
1,868,607✔
1320
  keyLen += nullFlagSize;
1,868,607✔
1321

1322
  keyBuf = taosMemoryCalloc(1, keyLen);
1,868,607✔
1323
  if (keyBuf == NULL) {
1,869,241✔
1324
    code = terrno;
×
UNCOV
1325
    goto end;
×
1326
  }
1327

1328
  if (initRemainGroups) {
1,869,241✔
1329
    pTableListInfo->remainGroups =
837,191✔
1330
        taosHashInit(rows, taosGetDefaultHashFunction(TSDB_DATA_TYPE_BIGINT), false, HASH_NO_LOCK);
836,201✔
1331
    if (pTableListInfo->remainGroups == NULL) {
837,191✔
1332
      code = terrno;
×
UNCOV
1333
      goto end;
×
1334
    }
1335
  }
1336

1337
  for (int i = 0; i < rows; i++) {
11,059,969✔
1338
    STableKeyInfo* info = taosArrayGet(pTableListInfo->pTableList, i);
9,189,991✔
1339
    QUERY_CHECK_NULL(info, code, lino, end, terrno);
9,191,016✔
1340

1341
    if (groupIdMap != NULL){
9,191,016✔
1342
      gInfo = taosArrayInit(taosArrayGetSize(groupData), sizeof(SStreamGroupValue));
157,772✔
1343
    }
1344
    
1345
    char* isNull = (char*)keyBuf;
9,191,071✔
1346
    char* pStart = (char*)keyBuf + sizeof(int8_t) * LIST_LENGTH(group);
9,191,071✔
1347
    for (int j = 0; j < taosArrayGetSize(groupData); j++) {
19,036,528✔
1348
      SColumnInfoData* pValue = (SColumnInfoData*)taosArrayGetP(groupData, j);
9,847,398✔
1349

1350
      if (groupIdMap != NULL && gInfo != NULL) {
9,847,398✔
1351
        int32_t ret = buildGroupInfo(pValue, i, gInfo);
180,347✔
1352
        if (ret != TSDB_CODE_SUCCESS) {
180,760✔
1353
          qError("buildGroupInfo failed at line %d since %s", __LINE__, tstrerror(ret));
×
1354
          taosArrayDestroyEx(gInfo, tDestroySStreamGroupValue);
×
UNCOV
1355
          gInfo = NULL;
×
1356
        }
1357
      }
1358
      
1359
      if (colDataIsNull_s(pValue, i)) {
19,695,468✔
1360
        isNull[j] = 1;
89,941✔
1361
      } else {
1362
        isNull[j] = 0;
9,757,716✔
1363
        char* data = colDataGetData(pValue, i);
9,756,871✔
1364
        if (pValue->info.type == TSDB_DATA_TYPE_JSON) {
9,758,019✔
1365
          // if (tTagIsJson(data)) {
1366
          //   code = TSDB_CODE_QRY_JSON_IN_GROUP_ERROR;
1367
          //   goto end;
1368
          // }
1369
          if (tTagIsJsonNull(data)) {
90,302✔
1370
            isNull[j] = 1;
×
UNCOV
1371
            continue;
×
1372
          }
1373
          int32_t len = getJsonValueLen(data);
90,302✔
1374
          memcpy(pStart, data, len);
90,302✔
1375
          pStart += len;
90,302✔
1376
        } else if (IS_VAR_DATA_TYPE(pValue->info.type)) {
9,666,570✔
1377
          if (IS_STR_DATA_BLOB(pValue->info.type)) {
6,632,909✔
1378
            if (blobDataTLen(data) > TSDB_MAX_BLOB_LEN) {
300✔
1379
              code = TSDB_CODE_TDB_INVALID_TABLE_SCHEMA_VER;
×
UNCOV
1380
              goto end;
×
1381
            }
1382
            memcpy(pStart, data, blobDataTLen(data));
×
UNCOV
1383
            pStart += blobDataTLen(data);
×
1384
          } else {
1385
            if (varDataTLen(data) > pValue->info.bytes) {
6,635,147✔
1386
              code = TSDB_CODE_TDB_INVALID_TABLE_SCHEMA_VER;
×
UNCOV
1387
              goto end;
×
1388
            }
1389
            memcpy(pStart, data, varDataTLen(data));
6,635,150✔
1390
            pStart += varDataTLen(data);
6,635,457✔
1391
          }
1392
        } else {
1393
          memcpy(pStart, data, pValue->info.bytes);
3,031,242✔
1394
          pStart += pValue->info.bytes;
3,031,937✔
1395
        }
1396
      }
1397
    }
1398

1399
    int32_t len = (int32_t)(pStart - (char*)keyBuf);
9,189,011✔
1400
    info->groupId = calcGroupId(keyBuf, len);
9,189,011✔
1401
    if (groupIdMap != NULL && gInfo != NULL) {
9,191,736✔
1402
      int32_t ret = taosHashPut(groupIdMap, &info->groupId, sizeof(info->groupId), &gInfo, POINTER_BYTES);
157,772✔
1403
      if (ret != TSDB_CODE_SUCCESS) {
157,772✔
1404
        qError("put groupid to map failed at line %d since %s", __LINE__, tstrerror(ret));
×
UNCOV
1405
        taosArrayDestroyEx(gInfo, tDestroySStreamGroupValue);
×
1406
      }
1407
      qDebug("put groupid to map gid:%" PRIu64, info->groupId);
157,772✔
1408
      gInfo = NULL;
157,772✔
1409
    }
1410
    if (initRemainGroups) {
9,191,736✔
1411
      // groupId ~ table uid
1412
      code = taosHashPut(pTableListInfo->remainGroups, &(info->groupId), sizeof(info->groupId), &(info->uid),
4,488,164✔
1413
                         sizeof(info->uid));
1414
      if (code == TSDB_CODE_DUP_KEY) {
4,486,486✔
1415
        code = TSDB_CODE_SUCCESS;
845,932✔
1416
      }
1417
      QUERY_CHECK_CODE(code, lino, end);
4,486,486✔
1418
    }
1419
  }
1420

1421
  if (tsTagFilterCache && groupIdMap == NULL) {
1,869,978✔
1422
    tableList = taosArrayDup(pTableListInfo->pTableList, NULL);
×
UNCOV
1423
    QUERY_CHECK_NULL(tableList, code, lino, end, terrno);
×
1424

UNCOV
1425
    code = pAPI->metaFn.metaPutTbGroupToCache(pVnode, pTableListInfo->idInfo.suid, context.digest,
×
1426
                                              tListLen(context.digest), tableList,
1427
                                              taosArrayGetSize(tableList) * sizeof(STableKeyInfo));
×
UNCOV
1428
    QUERY_CHECK_CODE(code, lino, end);
×
1429
  }
1430

1431
  //  int64_t st2 = taosGetTimestampUs();
1432
  //  qDebug("calculate tag block rows:%d, cost:%ld us", rows, st2-st1);
1433

1434
end:
1,869,828✔
1435
  taosMemoryFreeClear(keyBuf);
1,870,436✔
1436
  blockDataDestroy(pResBlock);
1,870,286✔
1437
  taosArrayDestroy(pBlockList);
1,869,536✔
1438
  taosArrayDestroyEx(pUidTagList, freeItem);
1,868,396✔
1439
  taosArrayDestroyP(groupData, releaseColInfoData);
1,870,286✔
1440
  taosArrayDestroyEx(gInfo, tDestroySStreamGroupValue);
1,870,136✔
1441

1442
  if (code != TSDB_CODE_SUCCESS) {
1,870,286✔
UNCOV
1443
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
1444
  }
1445
  return code;
1,868,901✔
1446
}
1447

1448
static int32_t nameComparFn(const void* p1, const void* p2) {
733,342✔
1449
  const char* pName1 = *(const char**)p1;
733,342✔
1450
  const char* pName2 = *(const char**)p2;
733,342✔
1451

1452
  int32_t ret = strcmp(pName1, pName2);
733,342✔
1453
  if (ret == 0) {
733,342✔
1454
    return 0;
18,342✔
1455
  } else {
1456
    return (ret > 0) ? 1 : -1;
715,000✔
1457
  }
1458
}
1459

1460
static SArray* getTableNameList(const SNodeListNode* pList) {
421,226✔
1461
  int32_t    code = TSDB_CODE_SUCCESS;
421,226✔
1462
  int32_t    lino = 0;
421,226✔
1463
  int32_t    len = LIST_LENGTH(pList->pNodeList);
421,226✔
1464
  SListCell* cell = pList->pNodeList->pHead;
421,226✔
1465

1466
  SArray* pTbList = taosArrayInit(len, POINTER_BYTES);
421,226✔
1467
  QUERY_CHECK_NULL(pTbList, code, lino, _end, terrno);
421,226✔
1468

1469
  for (int i = 0; i < pList->pNodeList->length; i++) {
1,144,705✔
1470
    SValueNode* valueNode = (SValueNode*)cell->pNode;
723,955✔
1471
    if (!IS_VAR_DATA_TYPE(valueNode->node.resType.type)) {
723,955✔
1472
      terrno = TSDB_CODE_INVALID_PARA;
×
1473
      taosArrayDestroy(pTbList);
×
UNCOV
1474
      return NULL;
×
1475
    }
1476

1477
    char* name = varDataVal(valueNode->datum.p);
723,955✔
1478
    void* tmp = taosArrayPush(pTbList, &name);
723,479✔
1479
    QUERY_CHECK_NULL(tmp, code, lino, _end, terrno);
723,479✔
1480
    cell = cell->pNext;
723,479✔
1481
  }
1482

1483
  size_t numOfTables = taosArrayGetSize(pTbList);
420,750✔
1484

1485
  // order the name
1486
  taosArraySort(pTbList, nameComparFn);
420,750✔
1487

1488
  // remove the duplicates
1489
  SArray* pNewList = taosArrayInit(taosArrayGetSize(pTbList), sizeof(void*));
421,226✔
1490
  QUERY_CHECK_NULL(pNewList, code, lino, _end, terrno);
421,226✔
1491
  void* tmpTbl = taosArrayGet(pTbList, 0);
421,226✔
1492
  QUERY_CHECK_NULL(tmpTbl, code, lino, _end, terrno);
421,226✔
1493
  void* tmp = taosArrayPush(pNewList, tmpTbl);
421,226✔
1494
  QUERY_CHECK_NULL(tmp, code, lino, _end, terrno);
421,226✔
1495

1496
  for (int32_t i = 1; i < numOfTables; ++i) {
723,955✔
1497
    char** name = taosArrayGetLast(pNewList);
302,729✔
1498
    char** nameInOldList = taosArrayGet(pTbList, i);
302,729✔
1499
    QUERY_CHECK_NULL(nameInOldList, code, lino, _end, terrno);
302,729✔
1500
    if (strcmp(*name, *nameInOldList) == 0) {
302,729✔
1501
      continue;
9,866✔
1502
    }
1503

1504
    tmp = taosArrayPush(pNewList, nameInOldList);
292,863✔
1505
    QUERY_CHECK_NULL(tmp, code, lino, _end, terrno);
292,863✔
1506
  }
1507

1508
_end:
421,226✔
1509
  taosArrayDestroy(pTbList);
421,226✔
1510
  if (code != TSDB_CODE_SUCCESS) {
421,226✔
1511
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
UNCOV
1512
    return NULL;
×
1513
  }
1514
  return pNewList;
421,226✔
1515
}
1516

1517
static int tableUidCompare(const void* a, const void* b) {
×
1518
  uint64_t u1 = *(uint64_t*)a;
×
UNCOV
1519
  uint64_t u2 = *(uint64_t*)b;
×
1520

1521
  if (u1 == u2) {
×
UNCOV
1522
    return 0;
×
1523
  }
1524

UNCOV
1525
  return u1 < u2 ? -1 : 1;
×
1526
}
1527

1528
static int32_t filterTableInfoCompare(const void* a, const void* b) {
18,488,667✔
1529
  STUidTagInfo* p1 = (STUidTagInfo*)a;
18,488,667✔
1530
  STUidTagInfo* p2 = (STUidTagInfo*)b;
18,488,667✔
1531

1532
  if (p1->uid == p2->uid) {
18,488,667✔
UNCOV
1533
    return 0;
×
1534
  }
1535

1536
  return p1->uid < p2->uid ? -1 : 1;
18,488,667✔
1537
}
1538

1539
static FilterCondType checkTagCond(SNode* cond) {
13,042,476✔
1540
  if (nodeType(cond) == QUERY_NODE_OPERATOR) {
13,042,476✔
1541
    return FILTER_NO_LOGIC;
11,103,435✔
1542
  }
1543
  if (nodeType(cond) == QUERY_NODE_LOGIC_CONDITION && ((SLogicConditionNode*)cond)->condType == LOGIC_COND_TYPE_AND) {
1,940,049✔
1544
    return FILTER_AND;
1,728,660✔
1545
  }
1546
  return FILTER_OTHER;
209,981✔
1547
}
1548

1549
static int32_t optimizeTbnameInCond(void* pVnode, int64_t suid, SArray* list, SNode* cond, SStorageAPI* pAPI) {
13,457,740✔
1550
  int32_t ret = -1;
13,457,740✔
1551
  int32_t ntype = nodeType(cond);
13,457,740✔
1552

1553
  if (ntype == QUERY_NODE_OPERATOR) {
13,458,358✔
1554
    ret = optimizeTbnameInCondImpl(pVnode, list, cond, pAPI, suid);
11,511,083✔
1555
    return ret;
11,510,784✔
1556
  }
1557
  if (ntype != QUERY_NODE_LOGIC_CONDITION || ((SLogicConditionNode*)cond)->condType != LOGIC_COND_TYPE_AND) {
1,947,275✔
1558
    return ret;
210,586✔
1559
  }
1560

1561
  bool                 hasTbnameCond = false;
1,735,889✔
1562
  SLogicConditionNode* pNode = (SLogicConditionNode*)cond;
1,735,889✔
1563
  SNodeList*           pList = (SNodeList*)pNode->pParameterList;
1,735,889✔
1564

1565
  int32_t len = LIST_LENGTH(pList);
1,735,684✔
1566
  if (len <= 0) {
1,736,489✔
UNCOV
1567
    return ret;
×
1568
  }
1569

1570
  SListCell* cell = pList->pHead;
1,736,489✔
1571
  for (int i = 0; i < len; i++) {
5,631,361✔
1572
    if (cell == NULL) break;
3,899,834✔
1573
    if (optimizeTbnameInCondImpl(pVnode, list, cell->pNode, pAPI, suid) == 0) {
3,899,834✔
1574
      hasTbnameCond = true;
6,294✔
1575
      break;
6,294✔
1576
    }
1577
    cell = cell->pNext;
3,894,459✔
1578
  }
1579

1580
  taosArraySort(list, filterTableInfoCompare);
1,737,821✔
1581
  taosArrayRemoveDuplicate(list, filterTableInfoCompare, NULL);
1,735,671✔
1582

1583
  if (hasTbnameCond) {
1,734,746✔
1584
    ret = pAPI->metaFn.getTableTagsByUid(pVnode, suid, list);
6,294✔
1585
  }
1586

1587
  return ret;
1,735,473✔
1588
}
1589

1590
// only return uid that does not contained in pExistedUidList
1591
static int32_t optimizeTbnameInCondImpl(void* pVnode, SArray* pExistedUidList, SNode* pTagCond, SStorageAPI* pStoreAPI,
15,413,459✔
1592
                                        uint64_t suid) {
1593
  if (nodeType(pTagCond) != QUERY_NODE_OPERATOR) {
15,413,459✔
1594
    return -1;
6,394✔
1595
  }
1596

1597
  SOperatorNode* pNode = (SOperatorNode*)pTagCond;
15,406,865✔
1598
  if (pNode->opType != OP_TYPE_IN) {
15,406,865✔
1599
    return -1;
14,677,654✔
1600
  }
1601

1602
  if ((pNode->pLeft != NULL && ((nodeType(pNode->pLeft) == QUERY_NODE_FUNCTION &&
729,206✔
1603
                                 ((SFunctionNode*)pNode->pLeft)->funcType == FUNCTION_TYPE_TBNAME)) ||
421,226✔
1604
       (nodeType(pNode->pLeft) == QUERY_NODE_COLUMN && ((SColumnNode*)pNode->pLeft)->colType == COLUMN_TYPE_TBNAME)) &&
307,980✔
1605
      (pNode->pRight != NULL && nodeType(pNode->pRight) == QUERY_NODE_NODE_LIST)) {
421,226✔
1606
    SNodeListNode* pList = (SNodeListNode*)pNode->pRight;
421,226✔
1607

1608
    int32_t len = LIST_LENGTH(pList->pNodeList);
421,226✔
1609
    if (len <= 0) {
421,226✔
UNCOV
1610
      return -1;
×
1611
    }
1612

1613
    SArray*   pTbList = getTableNameList(pList);
421,226✔
1614
    int32_t   numOfTables = taosArrayGetSize(pTbList);
421,226✔
1615
    SHashObj* uHash = NULL;
421,226✔
1616

1617
    size_t numOfExisted = taosArrayGetSize(pExistedUidList);  // len > 0 means there already have uids
421,226✔
1618
    if (numOfExisted > 0) {
421,226✔
1619
      uHash = taosHashInit(numOfExisted / 0.7, taosGetDefaultHashFunction(TSDB_DATA_TYPE_BIGINT), false, HASH_NO_LOCK);
2,362✔
1620
      if (!uHash) {
2,362✔
1621
        qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
UNCOV
1622
        return terrno;
×
1623
      }
1624

1625
      for (int i = 0; i < numOfExisted; i++) {
2,360,819✔
1626
        STUidTagInfo* pTInfo = taosArrayGet(pExistedUidList, i);
2,358,457✔
1627
        if (!pTInfo) {
2,358,457✔
1628
          qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
UNCOV
1629
          return terrno;
×
1630
        }
1631
        int32_t tempRes = taosHashPut(uHash, &pTInfo->uid, sizeof(uint64_t), &i, sizeof(i));
2,358,457✔
1632
        if (tempRes != TSDB_CODE_SUCCESS && tempRes != TSDB_CODE_DUP_KEY) {
2,358,457✔
1633
          qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(tempRes));
×
UNCOV
1634
          return tempRes;
×
1635
        }
1636
      }
1637
    }
1638

1639
    for (int i = 0; i < numOfTables; i++) {
1,118,887✔
1640
      char* name = taosArrayGetP(pTbList, i);
704,223✔
1641

1642
      uint64_t uid = 0, csuid = 0;
704,223✔
1643
      if (pStoreAPI->metaFn.getTableUidByName(pVnode, name, &uid) == 0) {
704,223✔
1644
        ETableType tbType = TSDB_TABLE_MAX;
421,348✔
1645
        if (pStoreAPI->metaFn.getTableTypeSuidByName(pVnode, name, &tbType, &csuid) == 0 &&
421,348✔
1646
            tbType == TSDB_CHILD_TABLE) {
421,348✔
1647
          if (suid != csuid) {
414,786✔
1648
            continue;
896✔
1649
          }
1650
          if (NULL == uHash || taosHashGet(uHash, &uid, sizeof(uid)) == NULL) {
413,890✔
1651
            STUidTagInfo s = {.uid = uid, .name = name, .pTagVal = NULL};
412,709✔
1652
            void*        tmp = taosArrayPush(pExistedUidList, &s);
412,709✔
1653
            if (!tmp) {
412,709✔
UNCOV
1654
              return terrno;
×
1655
            }
1656
          }
1657
        } else {
1658
          taosArrayDestroy(pTbList);
6,562✔
1659
          taosHashCleanup(uHash);
6,562✔
1660
          return -1;
6,562✔
1661
        }
1662
      } else {
1663
        //        qWarn("failed to get tableIds from by table name: %s, reason: %s", name, tstrerror(terrno));
1664
        terrno = 0;
282,875✔
1665
      }
1666
    }
1667

1668
    taosHashCleanup(uHash);
414,664✔
1669
    taosArrayDestroy(pTbList);
414,664✔
1670
    return 0;
414,664✔
1671
  }
1672

1673
  return -1;
307,980✔
1674
}
1675

1676
SSDataBlock* createTagValBlockForFilter(SArray* pColList, int32_t numOfTables, SArray* pUidTagList, void* pVnode,
14,235,795✔
1677
                                        SStorageAPI* pStorageAPI) {
1678
  int32_t      code = TSDB_CODE_SUCCESS;
14,235,795✔
1679
  int32_t      lino = 0;
14,235,795✔
1680
  SSDataBlock* pResBlock = NULL;
14,235,795✔
1681
  code = createDataBlock(&pResBlock);
14,236,163✔
1682
  QUERY_CHECK_CODE(code, lino, _end);
14,233,519✔
1683

1684
  for (int32_t i = 0; i < taosArrayGetSize(pColList); ++i) {
29,637,840✔
1685
    SColumnInfoData colInfo = {0};
15,403,409✔
1686
    void*           tmp = taosArrayGet(pColList, i);
15,402,098✔
1687
    QUERY_CHECK_NULL(tmp, code, lino, _end, terrno);
15,399,109✔
1688
    colInfo.info = *(SColumnInfo*)tmp;
15,399,109✔
1689
    code = blockDataAppendColInfo(pResBlock, &colInfo);
15,398,709✔
1690
    QUERY_CHECK_CODE(code, lino, _end);
15,403,886✔
1691
  }
1692

1693
  code = blockDataEnsureCapacity(pResBlock, numOfTables);
14,234,826✔
1694
  if (code != TSDB_CODE_SUCCESS) {
14,235,126✔
1695
    terrno = code;
×
1696
    blockDataDestroy(pResBlock);
×
UNCOV
1697
    return NULL;
×
1698
  }
1699

1700
  pResBlock->info.rows = numOfTables;
14,235,126✔
1701

1702
  int32_t numOfCols = taosArrayGetSize(pResBlock->pDataBlock);
14,237,135✔
1703

1704
  for (int32_t i = 0; i < numOfTables; i++) {
206,089,852✔
1705
    STUidTagInfo* p1 = taosArrayGet(pUidTagList, i);
191,851,519✔
1706
    QUERY_CHECK_NULL(p1, code, lino, _end, terrno);
191,852,333✔
1707

1708
    for (int32_t j = 0; j < numOfCols; j++) {
391,317,934✔
1709
      SColumnInfoData* pColInfo = (SColumnInfoData*)taosArrayGet(pResBlock->pDataBlock, j);
199,452,319✔
1710
      QUERY_CHECK_NULL(pColInfo, code, lino, _end, terrno);
199,430,604✔
1711

1712
      if (pColInfo->info.colId == -1) {  // tbname
199,430,604✔
1713
        char str[TSDB_TABLE_FNAME_LEN + VARSTR_HEADER_SIZE] = {0};
7,979,081✔
1714
        if (p1->name != NULL) {
7,979,671✔
1715
          STR_TO_VARSTR(str, p1->name);
412,709✔
1716
        } else {  // name is not retrieved during filter
1717
          code = pStorageAPI->metaFn.getTableNameByUid(pVnode, p1->uid, str);
7,566,754✔
1718
          QUERY_CHECK_CODE(code, lino, _end);
7,565,780✔
1719
        }
1720

1721
        code = colDataSetVal(pColInfo, i, str, false);
7,978,489✔
1722
        QUERY_CHECK_CODE(code, lino, _end);
7,976,718✔
1723
#if TAG_FILTER_DEBUG
1724
        qDebug("tagfilter uid:%ld, tbname:%s", *uid, str + 2);
1725
#endif
1726
      } else {
1727
        STagVal tagVal = {0};
191,460,150✔
1728
        tagVal.cid = pColInfo->info.colId;
191,460,846✔
1729
        if (p1->pTagVal == NULL) {
191,463,076✔
1730
          colDataSetNULL(pColInfo, i);
7,850✔
1731
        } else {
1732
          const char* p = pStorageAPI->metaFn.extractTagVal(p1->pTagVal, pColInfo->info.type, &tagVal);
191,445,171✔
1733

1734
          if (p == NULL || (pColInfo->info.type == TSDB_DATA_TYPE_JSON && ((STag*)p)->nTag == 0)) {
191,480,267✔
1735
            colDataSetNULL(pColInfo, i);
4,427,125✔
1736
          } else if (pColInfo->info.type == TSDB_DATA_TYPE_JSON) {
187,056,036✔
1737
            code = colDataSetVal(pColInfo, i, p, false);
714,250✔
1738
            QUERY_CHECK_CODE(code, lino, _end);
714,250✔
1739
          } else if (IS_VAR_DATA_TYPE(pColInfo->info.type)) {
301,281,965✔
1740
            if (IS_STR_DATA_BLOB(pColInfo->info.type)) {
114,947,734✔
UNCOV
1741
              QUERY_CHECK_CODE(code = TSDB_CODE_BLOB_NOT_SUPPORT_TAG, lino, _end);
×
1742
            }
1743
            char* tmp = taosMemoryMalloc(tagVal.nData + VARSTR_HEADER_SIZE + 1);
114,952,050✔
1744
            QUERY_CHECK_NULL(tmp, code, lino, _end, terrno);
114,942,787✔
1745
            varDataSetLen(tmp, tagVal.nData);
114,942,787✔
1746
            memcpy(tmp + VARSTR_HEADER_SIZE, tagVal.pData, tagVal.nData);
114,945,935✔
1747
            code = colDataSetVal(pColInfo, i, tmp, false);
114,944,104✔
1748
#if TAG_FILTER_DEBUG
1749
            qDebug("tagfilter varch:%s", tmp + 2);
1750
#endif
1751
            taosMemoryFree(tmp);
114,948,220✔
1752
            QUERY_CHECK_CODE(code, lino, _end);
114,947,258✔
1753
          } else {
1754
            code = colDataSetVal(pColInfo, i, (const char*)&tagVal.i64, false);
71,387,101✔
1755
            QUERY_CHECK_CODE(code, lino, _end);
71,391,994✔
1756
#if TAG_FILTER_DEBUG
1757
            if (pColInfo->info.type == TSDB_DATA_TYPE_INT) {
1758
              qDebug("tagfilter int:%d", *(int*)(&tagVal.i64));
1759
            } else if (pColInfo->info.type == TSDB_DATA_TYPE_DOUBLE) {
1760
              qDebug("tagfilter double:%f", *(double*)(&tagVal.i64));
1761
            }
1762
#endif
1763
          }
1764
        }
1765
      }
1766
    }
1767
  }
1768

1769
_end:
14,237,574✔
1770
  if (code != TSDB_CODE_SUCCESS) {
14,238,333✔
1771
    blockDataDestroy(pResBlock);
237✔
1772
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
1773
    terrno = code;
×
UNCOV
1774
    return NULL;
×
1775
  }
1776
  return pResBlock;
14,238,096✔
1777
}
1778

1779
static int32_t doSetQualifiedUid(STableListInfo* pListInfo, SArray* pUidList, const SArray* pUidTagList,
11,573,093✔
1780
                                 bool* pResultList, bool addUid) {
1781
  taosArrayClear(pUidList);
11,573,093✔
1782

1783
  STableKeyInfo info = {.uid = 0, .groupId = 0};
11,569,660✔
1784
  int32_t       numOfTables = taosArrayGetSize(pUidTagList);
11,570,068✔
1785
  for (int32_t i = 0; i < numOfTables; ++i) {
192,858,416✔
1786
    if (pResultList[i]) {
181,278,072✔
1787
      STUidTagInfo* tmpTag = (STUidTagInfo*)taosArrayGet(pUidTagList, i);
79,486,336✔
1788
      if (!tmpTag) {
79,484,152✔
1789
        qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
UNCOV
1790
        return terrno;
×
1791
      }
1792
      uint64_t uid = tmpTag->uid;
79,484,152✔
1793
      qDebug("tagfilter get uid:%" PRId64 ", res:%d", uid, pResultList[i]);
79,483,541✔
1794

1795
      info.uid = uid;
79,497,043✔
1796
      //qInfo("doSetQualifiedUid row:%d added to pTableList", i);
1797
      void* p = taosArrayPush(pListInfo->pTableList, &info);
79,497,043✔
1798
      if (p == NULL) {
79,495,771✔
UNCOV
1799
        return terrno;
×
1800
      }
1801

1802
      if (addUid) {
79,495,771✔
1803
        //qInfo("doSetQualifiedUid row:%d added to pUidList", i);
1804
        void* tmp = taosArrayPush(pUidList, &uid);
18,396✔
1805
        if (tmp == NULL) {
18,396✔
UNCOV
1806
          return terrno;
×
1807
        }
1808
      }
1809
    } else {
1810
      //qInfo("doSetQualifiedUid row:%d failed", i);
1811
    }
1812
  }
1813

1814
  return TSDB_CODE_SUCCESS;
11,580,344✔
1815
}
1816

1817
static int32_t copyExistedUids(SArray* pUidTagList, const SArray* pUidList) {
13,458,558✔
1818
  int32_t code = TSDB_CODE_SUCCESS;
13,458,558✔
1819
  int32_t numOfExisted = taosArrayGetSize(pUidList);
13,458,558✔
1820
  if (numOfExisted == 0) {
13,458,350✔
1821
    return code;
10,525,141✔
1822
  }
1823

1824
  for (int32_t i = 0; i < numOfExisted; ++i) {
35,922,039✔
1825
    uint64_t* uid = taosArrayGet(pUidList, i);
32,988,830✔
1826
    if (!uid) {
32,988,830✔
1827
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
UNCOV
1828
      return terrno;
×
1829
    }
1830
    STUidTagInfo info = {.uid = *uid};
32,988,830✔
1831
    void*        tmp = taosArrayPush(pUidTagList, &info);
32,988,364✔
1832
    if (!tmp) {
32,988,364✔
1833
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
UNCOV
1834
      return code;
×
1835
    }
1836
  }
1837
  return code;
2,933,209✔
1838
}
1839

1840
int32_t doFilterByTagCond(STableListInfo* pListInfo, SArray* pUidList, SNode* pTagCond, void* pVnode,
220,098,174✔
1841
                                 SIdxFltStatus status, SStorageAPI* pAPI, bool addUid, bool* listAdded, void* pStreamInfo) {
1842
  *listAdded = false;
220,098,174✔
1843
  if (pTagCond == NULL) {
220,128,343✔
1844
    return TSDB_CODE_SUCCESS;
206,623,390✔
1845
  }
1846

1847
  terrno = TSDB_CODE_SUCCESS;
13,504,953✔
1848

1849
  int32_t      lino = 0;
13,456,717✔
1850
  int32_t      code = TSDB_CODE_SUCCESS;
13,456,717✔
1851
  SArray*      pBlockList = NULL;
13,456,717✔
1852
  SSDataBlock* pResBlock = NULL;
13,456,717✔
1853
  SScalarParam output = {0};
13,456,309✔
1854
  SArray*      pUidTagList = NULL;
13,456,709✔
1855

1856
  SDataType type = {.type = TSDB_DATA_TYPE_BOOL, .bytes = sizeof(bool)};
13,456,709✔
1857

1858
  //  int64_t stt = taosGetTimestampUs();
1859
  pUidTagList = taosArrayInit(10, sizeof(STUidTagInfo));
13,456,937✔
1860
  QUERY_CHECK_NULL(pUidTagList, code, lino, end, terrno);
13,456,124✔
1861

1862
  code = copyExistedUids(pUidTagList, pUidList);
13,456,124✔
1863
  QUERY_CHECK_CODE(code, lino, end);
13,455,839✔
1864

1865
  int32_t filter = optimizeTbnameInCond(pVnode, pListInfo->idInfo.suid, pUidTagList, pTagCond, pAPI);
13,455,839✔
1866
  if (filter == 0) {  // tbname in filter is activated, do nothing and return
13,455,862✔
1867
    taosArrayClear(pUidList);
414,664✔
1868

1869
    int32_t numOfRows = taosArrayGetSize(pUidTagList);
414,664✔
1870
    code = taosArrayEnsureCap(pUidList, numOfRows);
414,664✔
1871
    QUERY_CHECK_CODE(code, lino, end);
414,664✔
1872

1873
    for (int32_t i = 0; i < numOfRows; ++i) {
3,185,830✔
1874
      STUidTagInfo* pInfo = taosArrayGet(pUidTagList, i);
2,771,166✔
1875
      QUERY_CHECK_NULL(pInfo, code, lino, end, terrno);
2,771,166✔
1876
      void* tmp = taosArrayPush(pUidList, &pInfo->uid);
2,771,166✔
1877
      QUERY_CHECK_NULL(tmp, code, lino, end, terrno);
2,771,166✔
1878
    }
1879
    terrno = 0;
414,664✔
1880
  } else {
1881
    qDebug("pUidTagList size:%d", (int32_t)taosArrayGetSize(pUidTagList));
13,041,198✔
1882

1883
    FilterCondType condType = checkTagCond(pTagCond);
13,041,456✔
1884
    if (((condType == FILTER_NO_LOGIC || condType == FILTER_AND) && status != SFLT_NOT_INDEX) ||
22,704,954✔
1885
          taosArrayGetSize(pUidTagList) > 0) {
9,662,075✔
1886
      code = pAPI->metaFn.getTableTagsByUid(pVnode, pListInfo->idInfo.suid, pUidTagList);
3,668,051✔
1887
    } else {
1888
      code = pAPI->metaFn.getTableTags(pVnode, pListInfo->idInfo.suid, pUidTagList);
9,374,828✔
1889
    }
1890
    if (code != TSDB_CODE_SUCCESS) {
13,042,852✔
1891
      qError("failed to get table tags from meta, reason:%s, suid:%" PRIu64, tstrerror(code), pListInfo->idInfo.suid);
×
1892
      terrno = code;
×
1893
      QUERY_CHECK_CODE(code, lino, end);
×
1894
    }
1895
  }
1896

1897
  qDebug("final pUidTagList size:%d", (int32_t)taosArrayGetSize(pUidTagList));
13,457,516✔
1898

1899
  int32_t numOfTables = taosArrayGetSize(pUidTagList);
13,458,628✔
1900
  if (numOfTables == 0) {
13,458,758✔
1901
    goto end;
1,875,808✔
1902
  }
1903

1904
  SArray* pColList = NULL;
11,582,950✔
1905
  code = qGetColumnsFromNodeList(pTagCond, false, &pColList); 
11,582,537✔
1906
  if (code != TSDB_CODE_SUCCESS) {
11,579,438✔
1907
    goto end;
×
1908
  }
1909
  pResBlock = createTagValBlockForFilter(pColList, numOfTables, pUidTagList, pVnode, pAPI);
11,579,438✔
1910
  taosArrayDestroy(pColList);
11,582,830✔
1911
  if (pResBlock == NULL) {
11,580,623✔
1912
    code = terrno;
×
1913
    QUERY_CHECK_CODE(code, lino, end);
×
1914
  }
1915

1916
  //fprintDataBlock(pResBlock, "tagFilter", "", 0);
1917

1918
  //  int64_t st1 = taosGetTimestampUs();
1919
  //  qDebug("generate tag block rows:%d, cost:%ld us", rows, st1-st);
1920
  pBlockList = taosArrayInit(2, POINTER_BYTES);
11,580,623✔
1921
  QUERY_CHECK_NULL(pBlockList, code, lino, end, terrno);
11,581,101✔
1922

1923
  void* tmp = taosArrayPush(pBlockList, &pResBlock);
11,582,308✔
1924
  QUERY_CHECK_NULL(tmp, code, lino, end, terrno);
11,582,308✔
1925

1926
  code = createResultData(&type, numOfTables, &output);
11,582,308✔
1927
  if (code != TSDB_CODE_SUCCESS) {
11,575,652✔
1928
    terrno = code;
×
1929
    QUERY_CHECK_CODE(code, lino, end);
×
1930
  }
1931

1932
  gTaskScalarExtra.pStreamInfo = pStreamInfo;
11,575,652✔
1933
  gTaskScalarExtra.pStreamRange = NULL;
11,575,652✔
1934
  code = scalarCalculate(pTagCond, pBlockList, &output, &gTaskScalarExtra);
11,578,270✔
1935
  if (code != TSDB_CODE_SUCCESS) {
11,576,319✔
1936
    qError("failed to calculate scalar, reason:%s", tstrerror(code));
1,108✔
1937
    terrno = code;
1,108✔
1938
    QUERY_CHECK_CODE(code, lino, end);
1,108✔
1939
  }
1940

1941
  code = doSetQualifiedUid(pListInfo, pUidList, pUidTagList, (bool*)output.columnData->pData, addUid);
11,575,211✔
1942
  if (code != TSDB_CODE_SUCCESS) {
11,580,744✔
1943
    terrno = code;
×
1944
    QUERY_CHECK_CODE(code, lino, end);
×
1945
  }
1946
  *listAdded = true;
11,580,744✔
1947

1948
end:
13,456,117✔
1949
  if (code != TSDB_CODE_SUCCESS) {
13,458,149✔
1950
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
1,108✔
1951
  }
1952
  blockDataDestroy(pResBlock);
13,458,149✔
1953
  taosArrayDestroy(pBlockList);
13,454,611✔
1954
  taosArrayDestroyEx(pUidTagList, freeItem);
13,450,813✔
1955

1956
  colDataDestroy(output.columnData);
13,455,792✔
1957
  taosMemoryFreeClear(output.columnData);
13,455,792✔
1958
  return code;
13,457,244✔
1959
}
1960

1961
typedef struct {
1962
  int32_t code;
1963
  SStreamRuntimeFuncInfo* pStreamRuntimeInfo;
1964
} PlaceHolderContext;
1965

1966
static EDealRes replacePlaceHolderColumn(SNode** pNode, void* pContext) {
127,694✔
1967
  PlaceHolderContext* pData = (PlaceHolderContext*)pContext;
127,694✔
1968
  if (QUERY_NODE_FUNCTION != nodeType((*pNode))) {
127,694✔
1969
    return DEAL_RES_CONTINUE;
104,760✔
1970
  }
1971
  SFunctionNode* pFuncNode = *(SFunctionNode**)(pNode);
22,934✔
1972
  if (!fmIsStreamPesudoColVal(pFuncNode->funcId)) {
22,934✔
1973
    return DEAL_RES_CONTINUE;
884✔
1974
  }
1975
  pData->code = fmSetStreamPseudoFuncParamVal(pFuncNode->funcId, pFuncNode->pParameterList, pData->pStreamRuntimeInfo);
22,050✔
1976
  if (pData->code != TSDB_CODE_SUCCESS) {
22,050✔
1977
    return DEAL_RES_ERROR;
×
1978
  }
1979
  SNode* pFirstParam = nodesListGetNode(pFuncNode->pParameterList, 0);
22,050✔
1980
  ((SValueNode*)pFirstParam)->translate = true;
22,050✔
1981
  SValueNode* res = NULL;
22,050✔
1982
  pData->code = nodesCloneNode(pFirstParam, (SNode**)&res);
22,050✔
1983
  if (NULL == res) {
22,050✔
1984
    return DEAL_RES_ERROR;
×
1985
  }
1986
  nodesDestroyNode(*pNode);
22,050✔
1987
  *pNode = (SNode*)res;
22,050✔
1988

1989
  return DEAL_RES_CONTINUE;
22,050✔
1990
}
1991

1992
static void extractTagColId(SOperatorNode* pOpNode, SArray* pColIdArray) {
34,540✔
1993
  SNode* pLeft = pOpNode->pLeft;
34,540✔
1994
  SNode* pRight = pOpNode->pRight;
34,540✔
1995
  SColumnNode* pColNode = nodeType(pLeft) == QUERY_NODE_COLUMN ?
34,540✔
1996
    (SColumnNode*)pLeft : (SColumnNode*)pRight;
34,540✔
1997

1998
  col_id_t colId = pColNode->colId;
34,540✔
1999
  void* _tmp = taosArrayPush(pColIdArray, &colId);
34,540✔
2000
}
34,540✔
2001

2002
static int32_t buildTagCondKey(
17,270✔
2003
  const SNode* pTagCond, char** pTagCondKey,
2004
  int32_t* tagCondKeyLen, SArray** pTagColIds) {
2005
  if (NULL == pTagCond ||
17,270✔
2006
    (nodeType(pTagCond) != QUERY_NODE_OPERATOR &&
17,270✔
2007
      nodeType(pTagCond) != QUERY_NODE_LOGIC_CONDITION)) {
17,270✔
2008
    qError("invalid parameter to extract tag filter symbol");
×
2009
    return TSDB_CODE_INTERNAL_ERROR;
×
2010
  }
2011
  int32_t code = TSDB_CODE_SUCCESS;
17,270✔
2012
  int32_t lino = 0;
17,270✔
2013
  *pTagColIds = taosArrayInit(4, sizeof(col_id_t));
17,270✔
2014

2015
  if (nodeType(pTagCond) == QUERY_NODE_OPERATOR) {
17,270✔
2016
    extractTagColId((SOperatorNode*)pTagCond, *pTagColIds);
×
2017
  } else if (nodeType(pTagCond) == QUERY_NODE_LOGIC_CONDITION) {
17,270✔
2018
    SNode* pChild = NULL;
17,270✔
2019
    FOREACH(pChild, ((SLogicConditionNode*)pTagCond)->pParameterList) {
51,810✔
2020
      extractTagColId((SOperatorNode*)pChild, *pTagColIds);
34,540✔
2021
    }
2022
  }
2023

2024
  taosArraySort(*pTagColIds, compareUint16Val);
17,270✔
2025

2026
  // encode ordered colIds into key string, separated by ','
2027
  *tagCondKeyLen =
34,540✔
2028
    (int32_t)(taosArrayGetSize(*pTagColIds) * (sizeof(col_id_t) + 1) - 1);
17,270✔
2029
  *pTagCondKey = (char*)taosMemoryCalloc(1, *tagCondKeyLen);
17,270✔
2030
  TSDB_CHECK_NULL(*pTagCondKey, code, lino, _end, terrno);
17,270✔
2031
  char* pStart = *pTagCondKey;
17,270✔
2032
  for (int32_t i = 0; i < taosArrayGetSize(*pTagColIds); ++i) {
51,810✔
2033
    col_id_t* pColId = (col_id_t*)taosArrayGet(*pTagColIds, i);
34,540✔
2034
    TSDB_CHECK_NULL(pColId, code, lino, _end, terrno);
34,540✔
2035
    memcpy(pStart, pColId, sizeof(col_id_t));
34,540✔
2036
    pStart += sizeof(col_id_t);
34,540✔
2037
    if (i != taosArrayGetSize(*pTagColIds) - 1) {
34,540✔
2038
      *pStart = ',';
17,270✔
2039
      pStart += 1;
17,270✔
2040
    }
2041
  }
2042

2043
_end:
17,270✔
2044
  if (TSDB_CODE_SUCCESS != code) {
17,270✔
2045
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
2046
    terrno = code;
×
2047
  }
2048
  return code;
17,270✔
2049
}
2050

2051
static EDealRes canOptimizeTagCondFilter(SNode* pTagCond, void* pContext) {
161,396✔
2052
  if (NULL == pTagCond) {
161,396✔
2053
    *(bool*)pContext = false;
×
2054
    return DEAL_RES_END;
×
2055
  }
2056
  if (nodeType(pTagCond) == QUERY_NODE_VALUE ||
161,396✔
2057
    nodeType(pTagCond) == QUERY_NODE_COLUMN) {
107,074✔
2058
    return DEAL_RES_CONTINUE;
88,862✔
2059
  }
2060
  if (nodeType(pTagCond) == QUERY_NODE_OPERATOR &&
72,534✔
2061
    ((SOperatorNode*)pTagCond)->opType == OP_TYPE_EQUAL) {
35,482✔
2062
    return DEAL_RES_CONTINUE;
34,540✔
2063
  }
2064
  if (nodeType(pTagCond) == QUERY_NODE_LOGIC_CONDITION &&
37,994✔
2065
    ((SLogicConditionNode*)pTagCond)->condType == LOGIC_COND_TYPE_AND) {
17,270✔
2066
    return DEAL_RES_CONTINUE;
17,270✔
2067
  }
2068
  if (nodeType(pTagCond) == QUERY_NODE_FUNCTION &&
40,506✔
2069
    fmIsStreamPesudoColVal(((SFunctionNode*)pTagCond)->funcId)) {
19,782✔
2070
    return DEAL_RES_CONTINUE;
19,782✔
2071
  }
2072
  *(bool*)pContext = false;
942✔
2073
  return DEAL_RES_END;
942✔
2074
}
2075

2076
int32_t getTableList(void* pVnode, SScanPhysiNode* pScanNode, SNode* pTagCond, SNode* pTagIndexCond,
219,976,579✔
2077
                     STableListInfo* pListInfo, uint8_t* digest, const char* idstr, SStorageAPI* pStorageAPI, void* pStreamInfo) {
2078
  int32_t code = TSDB_CODE_SUCCESS;
219,976,579✔
2079
  int32_t lino = 0;
219,976,579✔
2080
  size_t  numOfTables = 0;
219,976,579✔
2081
  bool    listAdded = false;
219,976,579✔
2082

2083
  pListInfo->idInfo.suid = pScanNode->suid;
220,022,268✔
2084
  pListInfo->idInfo.tableType = pScanNode->tableType;
219,914,934✔
2085

2086
  SArray* pUidList = taosArrayInit(8, sizeof(uint64_t));
219,990,437✔
2087
  QUERY_CHECK_NULL(pUidList, code, lino, _error, terrno);
219,772,963✔
2088

2089
  SIdxFltStatus status = SFLT_NOT_INDEX;
219,772,963✔
2090
  char*   pTagCondKey = NULL;
219,793,980✔
2091
  int32_t tagCondKeyLen;
219,806,847✔
2092
  SArray* pTagColIds = NULL;
219,917,422✔
2093
  char*   pPayload = NULL;
219,947,750✔
2094
  qTrace("getTableList called, suid:%" PRIu64
219,947,750✔
2095
    ", tagCond:%p, tagIndexCond:%p, %d %d", pScanNode->suid, pTagCond,
2096
    pTagIndexCond, pScanNode->tableType, pScanNode->virtualStableScan);
2097
  if (pScanNode->tableType != TSDB_SUPER_TABLE && !pScanNode->virtualStableScan) {
219,947,750✔
2098
    pListInfo->idInfo.uid = pScanNode->uid;
150,628,614✔
2099
    if (pStorageAPI->metaFn.isTableExisted(pVnode, pScanNode->uid)) {
150,519,721✔
2100
      void* tmp = taosArrayPush(pUidList, &pScanNode->uid);
150,663,864✔
2101
      QUERY_CHECK_NULL(tmp, code, lino, _error, terrno);
150,669,608✔
2102
    }
2103
    code = doFilterByTagCond(pListInfo, pUidList, pTagCond, pVnode, status, pStorageAPI, false, &listAdded, pStreamInfo);
150,710,299✔
2104
    QUERY_CHECK_CODE(code, lino, _end);
150,705,015✔
2105
  } else {
2106
    bool      isStream = (pStreamInfo != NULL);
69,378,063✔
2107
    bool      hasTagCond = (pTagCond != NULL);
69,378,063✔
2108
    bool      canCacheTagEqCondFilter = false;
69,378,063✔
2109
    T_MD5_CTX context = {0};
69,336,996✔
2110

2111
    qTrace("start to get table list by tag filter, suid:%" PRIu64
69,392,050✔
2112
      ",tsStableTagFilterCache:%d, tsTagFilterCache:%d", 
2113
      pScanNode->suid, tsStableTagFilterCache, tsTagFilterCache);
2114

2115
    bool acquired = false;
69,392,050✔
2116
    // first, check whether we can use stable tag filter cache
2117
    if (tsStableTagFilterCache && isStream && hasTagCond) {
69,333,906✔
2118
      canCacheTagEqCondFilter = true;
18,212✔
2119
      nodesWalkExpr(pTagCond, canOptimizeTagCondFilter,
18,212✔
2120
        (void*)&canCacheTagEqCondFilter);
2121
    }
2122
    if (canCacheTagEqCondFilter) {
69,181,914✔
2123
      qDebug("%s, stable tag filter condition can be optimized", idstr);
17,270✔
2124
      if (((SStreamRuntimeFuncInfo*)pStreamInfo)->hasPlaceHolder) {
17,270✔
2125
        SNode* tmp = NULL;
17,270✔
2126
        code = nodesCloneNode((SNode*)pTagCond, &tmp);
17,270✔
2127
        QUERY_CHECK_CODE(code, lino, _error);
17,270✔
2128

2129
        PlaceHolderContext ctx = {.code = TSDB_CODE_SUCCESS, .pStreamRuntimeInfo = (SStreamRuntimeFuncInfo*)pStreamInfo};
17,270✔
2130
        nodesRewriteExpr(&tmp, replacePlaceHolderColumn, (void*)&ctx);
17,270✔
2131
        if (TSDB_CODE_SUCCESS != ctx.code) {
17,270✔
2132
          nodesDestroyNode(tmp);
×
2133
          code = ctx.code;
×
2134
          goto _error;
×
2135
        }
2136
        code = genStableTagFilterDigest(tmp, &context);
17,270✔
2137
        nodesDestroyNode(tmp);
17,270✔
2138
      } else {
2139
        code = genStableTagFilterDigest(pTagCond, &context);
×
2140
      }
2141
      QUERY_CHECK_CODE(code, lino, _error);
17,270✔
2142

2143
      code = buildTagCondKey(
17,270✔
2144
        pTagCond, &pTagCondKey, &tagCondKeyLen, &pTagColIds);
2145
      QUERY_CHECK_CODE(code, lino, _error);
17,270✔
2146
      code = pStorageAPI->metaFn.getStableCachedTableList(
17,270✔
2147
        pVnode, pScanNode->suid, pTagCondKey, tagCondKeyLen,
17,270✔
2148
        context.digest, tListLen(context.digest), pUidList, &acquired);
2149
      QUERY_CHECK_CODE(code, lino, _error);
17,270✔
2150
    } else if (tsTagFilterCache) {
69,164,644✔
2151
      // second, try to use normal tag filter cache
2152
      qDebug("%s using normal tag filter cache", idstr);
58,138✔
2153
      if (pStreamInfo != NULL && ((SStreamRuntimeFuncInfo*)pStreamInfo)->hasPlaceHolder) {
60,406✔
2154
        SNode* tmp = NULL;
2,268✔
2155
        code = nodesCloneNode((SNode*)pTagCond, &tmp);
2,268✔
2156
        QUERY_CHECK_CODE(code, lino, _error);
2,268✔
2157

2158
        PlaceHolderContext ctx = {.code = TSDB_CODE_SUCCESS, .pStreamRuntimeInfo = (SStreamRuntimeFuncInfo*)pStreamInfo};
2,268✔
2159
        nodesRewriteExpr(&tmp, replacePlaceHolderColumn, (void*)&ctx);
2,268✔
2160
        if (TSDB_CODE_SUCCESS != ctx.code) {
2,268✔
2161
          nodesDestroyNode(tmp);
×
2162
          code = ctx.code;
×
2163
          goto _error;
×
2164
        }
2165
        code = genTagFilterDigest(tmp, &context);
2,268✔
2166
        nodesDestroyNode(tmp);
2,268✔
2167
      } else {
2168
        code = genTagFilterDigest(pTagCond, &context);
55,870✔
2169
      }
2170
      // try to retrieve the result from meta cache
2171
      QUERY_CHECK_CODE(code, lino, _error);      
58,138✔
2172
      code = pStorageAPI->metaFn.getCachedTableList(
58,138✔
2173
        pVnode, pScanNode->suid, context.digest,
58,138✔
2174
        tListLen(context.digest), pUidList, &acquired);
2175
      QUERY_CHECK_CODE(code, lino, _error);
175,976✔
2176
    }
2177
    if (acquired) {
69,299,723✔
2178
      taosArrayDestroy(pTagColIds);
51,334✔
2179
      pTagColIds = NULL;
51,334✔
2180
      
2181
      digest[0] = 1;
51,334✔
2182
      memcpy(
102,668✔
2183
        digest + 1, context.digest, tListLen(context.digest));
51,334✔
2184
      qDebug("suid:%" PRIu64 ", %s retrieve table uid list from cache,"
51,334✔
2185
        " numOfTables:%d", 
2186
        pScanNode->suid, idstr, (int32_t)taosArrayGetSize(pUidList));
2187
      goto _end;
51,334✔
2188
    } else {
2189
      qDebug("suid:%" PRIu64 
69,248,389✔
2190
        ", failed to get table uid list from cache", pScanNode->suid);
2191
    }
2192

2193
    if (!pTagCond) {  // no tag filter condition exists, let's fetch all tables of this super table
69,327,292✔
2194
      code = pStorageAPI->metaFn.getChildTableList(pVnode, pScanNode->suid, pUidList);
56,146,980✔
2195
      QUERY_CHECK_CODE(code, lino, _error);
56,129,123✔
2196
      qTrace("no tag filter, get all child tables, numOfTables:%d", (int32_t)taosArrayGetSize(pUidList));
56,129,123✔
2197
    } else {
2198
      // failed to find the result in the cache, let try to calculate the results
2199
      if (pTagIndexCond) {
13,180,312✔
2200
        void* pIndex = pStorageAPI->metaFn.getInvertIndex(pVnode);
4,476,532✔
2201

2202
        SIndexMetaArg metaArg = {.metaEx = pVnode,
4,476,596✔
2203
                                 .idx = pStorageAPI->metaFn.storeGetIndexInfo(pVnode),
4,476,532✔
2204
                                 .ivtIdx = pIndex,
2205
                                 .suid = pScanNode->uid};
4,476,532✔
2206

2207
        status = SFLT_NOT_INDEX;
4,476,532✔
2208
        code = doFilterTag(pTagIndexCond, &metaArg, pUidList, &status, &pStorageAPI->metaFilter);
4,476,532✔
2209
        if (code != 0 || status == SFLT_NOT_INDEX) {  // temporarily disable it for performance sake
4,470,459✔
2210
          qDebug("failed to get tableIds from index, suid:%" PRIu64 ", uidListSize:%d", pScanNode->uid, (int32_t)taosArrayGetSize(pUidList));
1,084,088✔
2211
        } else {
2212
          qDebug("succ to get filter result, table num: %d", (int)taosArrayGetSize(pUidList));
3,386,371✔
2213
        }
2214
      }
2215
    }
2216
    qTrace("after index filter, pTagCond:%p uidListSize:%d", pTagCond, (int32_t)taosArrayGetSize(pUidList));
69,312,587✔
2217
    code = doFilterByTagCond(pListInfo, pUidList, pTagCond, pVnode, status,
69,348,831✔
2218
      pStorageAPI, tsTagFilterCache || tsStableTagFilterCache,
69,348,831✔
2219
      &listAdded, pStreamInfo);
2220
    QUERY_CHECK_CODE(code, lino, _error);
69,326,273✔
2221

2222
    // let's add the filter results into meta-cache
2223
    numOfTables = taosArrayGetSize(pUidList);
69,325,165✔
2224

2225
    if (canCacheTagEqCondFilter) {
69,326,064✔
2226
      qInfo("%s, suid:%" PRIu64 ", add uid list to stable tag filter cache, "
8,792✔
2227
            "uidListSize:%d, origin key:%" PRIu64 ",%" PRIu64,
2228
            idstr, pScanNode->suid, (int32_t)numOfTables,
2229
            *(uint64_t*)context.digest, *(uint64_t*)(context.digest + 8));
2230

2231
      code = pStorageAPI->metaFn.putStableCachedTableList(
8,792✔
2232
        pVnode, pScanNode->suid, pTagCondKey, tagCondKeyLen,
2233
        context.digest, tListLen(context.digest),
2234
        pUidList, &pTagColIds);
2235
      QUERY_CHECK_CODE(code, lino, _end);
8,792✔
2236

2237
      digest[0] = 1;
8,792✔
2238
      memcpy(digest + 1, context.digest, tListLen(context.digest));
8,792✔
2239
    } else if (tsTagFilterCache) {
69,317,272✔
2240
      qInfo("%s, suid:%" PRIu64 ", add uid list to normal tag filter cache, "
15,282✔
2241
            "uidListSize:%d, origin key:%" PRIu64 ",%" PRIu64,
2242
            idstr, pScanNode->suid, (int32_t)numOfTables,
2243
            *(uint64_t*)context.digest, *(uint64_t*)(context.digest + 8));
2244
      size_t size = numOfTables * sizeof(uint64_t) + sizeof(int32_t);
15,282✔
2245
      pPayload = taosMemoryMalloc(size);
15,282✔
2246
      QUERY_CHECK_NULL(pPayload, code, lino, _end, terrno);
15,282✔
2247

2248
      *(int32_t*)pPayload = (int32_t)numOfTables;
15,282✔
2249
      if (numOfTables > 0) {
15,282✔
2250
        void* tmp = taosArrayGet(pUidList, 0);
12,200✔
2251
        QUERY_CHECK_NULL(tmp, code, lino, _end, terrno);
12,200✔
2252
        memcpy(pPayload + sizeof(int32_t), tmp, numOfTables * sizeof(uint64_t));
12,200✔
2253
      }
2254

2255
      code = pStorageAPI->metaFn.putCachedTableList(pVnode, pScanNode->suid,
15,282✔
2256
                                                    context.digest,
2257
                                                    tListLen(context.digest),
2258
                                                    pPayload, size, 1);
2259
      if (TSDB_CODE_SUCCESS == code) {
15,282✔
2260
        /*
2261
          data referenced by pPayload is used in lru cache,
2262
          reset pPayload to NULL to avoid being freed in _error block
2263
        */
2264
        pPayload = NULL;
14,654✔
2265
      } else {
2266
        if (TSDB_CODE_DUP_KEY == code) {
628✔
2267
          /*
2268
            another thread has already put the same key into cache,
2269
            we can just ignore this error
2270
          */
2271
          code = TSDB_CODE_SUCCESS;
628✔
2272
        }
2273
        QUERY_CHECK_CODE(code, lino, _end);
628✔
2274
      }
2275

2276

2277
      digest[0] = 1;
15,282✔
2278
      memcpy(digest + 1, context.digest, tListLen(context.digest));
15,282✔
2279
    }
2280
  }
2281

2282
_end:
220,072,278✔
2283
  if (!listAdded) {
220,028,538✔
2284
    numOfTables = taosArrayGetSize(pUidList);
208,450,064✔
2285
    for (int i = 0; i < numOfTables; i++) {
628,804,298✔
2286
      void* tmp = taosArrayGet(pUidList, i);
420,316,861✔
2287
      QUERY_CHECK_NULL(tmp, code, lino, _error, terrno);
420,351,275✔
2288
      STableKeyInfo info = {.uid = *(uint64_t*)tmp, .groupId = 0};
420,351,275✔
2289

2290
      void* p = taosArrayPush(pListInfo->pTableList, &info);
420,320,726✔
2291
      if (p == NULL) {
420,388,063✔
2292
        taosArrayDestroy(pUidList);
×
2293
        return terrno;
×
2294
      }
2295

2296
      qTrace("tagfilter get uid:%" PRIu64 ", %s", info.uid, idstr);
420,388,063✔
2297
    }
2298
  }
2299

2300
  qDebug("%s, table list with %d uids built", idstr, (int32_t)numOfTables);
220,065,911✔
2301

2302
_error:
220,075,276✔
2303
  taosArrayDestroy(pUidList);
220,098,616✔
2304
  taosArrayDestroy(pTagColIds);
220,002,663✔
2305
  taosMemFreeClear(pTagCondKey);
220,026,156✔
2306
  taosMemFreeClear(pPayload);
220,026,156✔
2307
  if (code != TSDB_CODE_SUCCESS) {
220,026,156✔
2308
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
1,108✔
2309
  }
2310
  return code;
220,011,795✔
2311
}
2312

2313
int32_t qGetTableList(int64_t suid, void* pVnode, void* node, SArray** tableList, void* pTaskInfo) {
4,303✔
2314
  int32_t        code = TSDB_CODE_SUCCESS;
4,303✔
2315
  int32_t        lino = 0;
4,303✔
2316
  SSubplan*      pSubplan = (SSubplan*)node;
4,303✔
2317
  SScanPhysiNode pNode = {0};
4,303✔
2318
  pNode.suid = suid;
4,303✔
2319
  pNode.uid = suid;
4,303✔
2320
  pNode.tableType = TSDB_SUPER_TABLE;
4,303✔
2321

2322
  STableListInfo* pTableListInfo = tableListCreate();
4,303✔
2323
  QUERY_CHECK_NULL(pTableListInfo, code, lino, _end, terrno);
4,303✔
2324
  uint8_t digest[17] = {0};
4,303✔
2325
  code = getTableList(pVnode, &pNode, pSubplan ? pSubplan->pTagCond : NULL, pSubplan ? pSubplan->pTagIndexCond : NULL,
4,303✔
2326
                      pTableListInfo, digest, "qGetTableList", &((SExecTaskInfo*)pTaskInfo)->storageAPI, NULL);
2327
  QUERY_CHECK_CODE(code, lino, _end);
4,303✔
2328
  *tableList = pTableListInfo->pTableList;
4,303✔
2329
  pTableListInfo->pTableList = NULL;
4,303✔
2330
  tableListDestroy(pTableListInfo);
4,303✔
2331

2332
_end:
4,303✔
2333
  if (code != TSDB_CODE_SUCCESS) {
4,303✔
2334
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
2335
  }
2336
  return code;
4,303✔
2337
}
2338

2339
size_t getTableTagsBufLen(const SNodeList* pGroups) {
×
2340
  size_t keyLen = 0;
×
2341

2342
  SNode* node;
2343
  FOREACH(node, pGroups) {
×
2344
    SExprNode* pExpr = (SExprNode*)node;
×
2345
    keyLen += pExpr->resType.bytes;
×
2346
  }
2347

2348
  keyLen += sizeof(int8_t) * LIST_LENGTH(pGroups);
×
2349
  return keyLen;
×
2350
}
2351

2352
int32_t getGroupIdFromTagsVal(void* pVnode, uint64_t uid, SNodeList* pGroupNode, char* keyBuf, uint64_t* pGroupId,
×
2353
                              SStorageAPI* pAPI) {
2354
  SMetaReader mr = {0};
×
2355

2356
  pAPI->metaReaderFn.initReader(&mr, pVnode, META_READER_LOCK, &pAPI->metaFn);
×
2357
  if (pAPI->metaReaderFn.getEntryGetUidCache(&mr, uid) != 0) {  // table not exist
×
2358
    pAPI->metaReaderFn.clearReader(&mr);
×
2359
    return TSDB_CODE_PAR_TABLE_NOT_EXIST;
×
2360
  }
2361

2362
  SNodeList* groupNew = NULL;
×
2363
  int32_t    code = nodesCloneList(pGroupNode, &groupNew);
×
2364
  if (TSDB_CODE_SUCCESS != code) {
×
2365
    pAPI->metaReaderFn.clearReader(&mr);
×
2366
    return code;
×
2367
  }
2368

2369
  STransTagExprCtx ctx = {.code = 0, .pReader = &mr};
×
2370
  nodesRewriteExprsPostOrder(groupNew, doTranslateTagExpr, &ctx);
×
2371
  if (TSDB_CODE_SUCCESS != ctx.code) {
×
2372
    nodesDestroyList(groupNew);
×
2373
    pAPI->metaReaderFn.clearReader(&mr);
×
2374
    return code;
×
2375
  }
2376
  char* isNull = (char*)keyBuf;
×
2377
  char* pStart = (char*)keyBuf + sizeof(int8_t) * LIST_LENGTH(pGroupNode);
×
2378

2379
  SNode*  pNode;
2380
  int32_t index = 0;
×
2381
  FOREACH(pNode, groupNew) {
×
2382
    SNode*  pNew = NULL;
×
2383
    int32_t code = scalarCalculateConstants(pNode, &pNew);
×
2384
    if (TSDB_CODE_SUCCESS == code) {
×
2385
      REPLACE_NODE(pNew);
×
2386
    } else {
2387
      nodesDestroyList(groupNew);
×
2388
      pAPI->metaReaderFn.clearReader(&mr);
×
2389
      return code;
×
2390
    }
2391

2392
    if (nodeType(pNew) != QUERY_NODE_VALUE) {
×
2393
      nodesDestroyList(groupNew);
×
2394
      pAPI->metaReaderFn.clearReader(&mr);
×
2395
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR));
×
2396
      return TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR;
×
2397
    }
2398
    SValueNode* pValue = (SValueNode*)pNew;
×
2399

2400
    if (pValue->node.resType.type == TSDB_DATA_TYPE_NULL || pValue->isNull) {
×
2401
      isNull[index++] = 1;
×
2402
      continue;
×
2403
    } else {
2404
      isNull[index++] = 0;
×
2405
      char* data = nodesGetValueFromNode(pValue);
×
2406
      if (pValue->node.resType.type == TSDB_DATA_TYPE_JSON) {
×
2407
        if (tTagIsJson(data)) {
×
2408
          terrno = TSDB_CODE_QRY_JSON_IN_GROUP_ERROR;
×
2409
          nodesDestroyList(groupNew);
×
2410
          pAPI->metaReaderFn.clearReader(&mr);
×
2411
          return terrno;
×
2412
        }
2413
        int32_t len = getJsonValueLen(data);
×
2414
        memcpy(pStart, data, len);
×
2415
        pStart += len;
×
2416
      } else if (IS_VAR_DATA_TYPE(pValue->node.resType.type)) {
×
2417
        if (IS_STR_DATA_BLOB(pValue->node.resType.type)) {
×
2418
          return TSDB_CODE_BLOB_NOT_SUPPORT_TAG;
×
2419
        }
2420
        memcpy(pStart, data, varDataTLen(data));
×
2421
        pStart += varDataTLen(data);
×
2422
      } else {
2423
        memcpy(pStart, data, pValue->node.resType.bytes);
×
2424
        pStart += pValue->node.resType.bytes;
×
2425
      }
2426
    }
2427
  }
2428

2429
  int32_t len = (int32_t)(pStart - (char*)keyBuf);
×
2430
  *pGroupId = calcGroupId(keyBuf, len);
×
2431

2432
  nodesDestroyList(groupNew);
×
2433
  pAPI->metaReaderFn.clearReader(&mr);
×
2434

2435
  return TSDB_CODE_SUCCESS;
×
2436
}
2437

2438
SArray* makeColumnArrayFromList(SNodeList* pNodeList) {
8,133,744✔
2439
  if (!pNodeList) {
8,133,744✔
2440
    return NULL;
×
2441
  }
2442

2443
  size_t  numOfCols = LIST_LENGTH(pNodeList);
8,133,744✔
2444
  SArray* pList = taosArrayInit(numOfCols, sizeof(SColumn));
8,133,744✔
2445
  if (pList == NULL) {
8,132,080✔
2446
    return NULL;
×
2447
  }
2448

2449
  for (int32_t i = 0; i < numOfCols; ++i) {
18,280,981✔
2450
    SColumnNode* pColNode = (SColumnNode*)nodesListGetNode(pNodeList, i);
10,148,549✔
2451
    if (!pColNode) {
10,149,962✔
2452
      taosArrayDestroy(pList);
×
2453
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR));
×
2454
      return NULL;
×
2455
    }
2456

2457
    // todo extract method
2458
    SColumn c = {0};
10,149,962✔
2459
    c.slotId = pColNode->slotId;
10,150,565✔
2460
    c.colId = pColNode->colId;
10,151,168✔
2461
    c.type = pColNode->node.resType.type;
10,150,565✔
2462
    c.bytes = pColNode->node.resType.bytes;
10,151,270✔
2463
    c.precision = pColNode->node.resType.precision;
10,151,023✔
2464
    c.scale = pColNode->node.resType.scale;
10,149,214✔
2465

2466
    void* tmp = taosArrayPush(pList, &c);
10,151,023✔
2467
    if (!tmp) {
10,151,023✔
2468
      taosArrayDestroy(pList);
×
2469
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
2470
      return NULL;
×
2471
    }
2472
  }
2473

2474
  return pList;
8,132,432✔
2475
}
2476

2477
int32_t extractColMatchInfo(SNodeList* pNodeList, SDataBlockDescNode* pOutputNodeList, int32_t* numOfOutputCols,
269,179,684✔
2478
                            int32_t type, SColMatchInfo* pMatchInfo) {
2479
  size_t  numOfCols = LIST_LENGTH(pNodeList);
269,179,684✔
2480
  int32_t code = TSDB_CODE_SUCCESS;
269,208,680✔
2481
  int32_t lino = 0;
269,208,680✔
2482

2483
  pMatchInfo->matchType = type;
269,208,680✔
2484

2485
  SArray* pList = taosArrayInit(numOfCols, sizeof(SColMatchItem));
269,035,033✔
2486
  if (pList == NULL) {
269,051,399✔
2487
    code = terrno;
×
2488
    return code;
×
2489
  }
2490

2491
  for (int32_t i = 0; i < numOfCols; ++i) {
1,231,439,094✔
2492
    STargetNode* pNode = (STargetNode*)nodesListGetNode(pNodeList, i);
962,222,071✔
2493
    QUERY_CHECK_NULL(pNode, code, lino, _end, terrno);
962,328,648✔
2494
    if (nodeType(pNode->pExpr) == QUERY_NODE_COLUMN) {
962,328,648✔
2495
      SColumnNode* pColNode = (SColumnNode*)pNode->pExpr;
957,959,446✔
2496

2497
      SColMatchItem c = {.needOutput = true};
957,959,765✔
2498
      c.colId = pColNode->colId;
957,979,441✔
2499
      c.srcSlotId = pColNode->slotId;
957,947,046✔
2500
      c.dstSlotId = pNode->slotId;
957,925,116✔
2501
      c.isPk = pColNode->isPk;
957,914,344✔
2502
      c.dataType = pColNode->node.resType;
957,886,119✔
2503
      void* tmp = taosArrayPush(pList, &c);
957,940,492✔
2504
      QUERY_CHECK_NULL(tmp, code, lino, _end, terrno);
957,940,492✔
2505
    }
2506
  }
2507

2508
  // set the output flag for each column in SColMatchInfo, according to the
2509
  *numOfOutputCols = 0;
269,217,023✔
2510
  int32_t num = LIST_LENGTH(pOutputNodeList->pSlots);
269,255,385✔
2511
  for (int32_t i = 0; i < num; ++i) {
1,333,598,238✔
2512
    SSlotDescNode* pNode = (SSlotDescNode*)nodesListGetNode(pOutputNodeList->pSlots, i);
1,064,377,349✔
2513
    QUERY_CHECK_NULL(pNode, code, lino, _end, terrno);
1,064,461,037✔
2514

2515
    // todo: add reserve flag check
2516
    // it is a column reserved for the arithmetic expression calculation
2517
    if (pNode->slotId >= numOfCols) {
1,064,461,037✔
2518
      (*numOfOutputCols) += 1;
102,096,946✔
2519
      continue;
102,098,647✔
2520
    }
2521

2522
    SColMatchItem* info = NULL;
962,387,610✔
2523
    for (int32_t j = 0; j < taosArrayGetSize(pList); ++j) {
2,147,483,647✔
2524
      info = taosArrayGet(pList, j);
2,147,483,647✔
2525
      QUERY_CHECK_NULL(info, code, lino, _end, terrno);
2,147,483,647✔
2526
      if (info->dstSlotId == pNode->slotId) {
2,147,483,647✔
2527
        break;
957,046,061✔
2528
      }
2529
    }
2530

2531
    if (pNode->output) {
9,797,757✔
2532
      (*numOfOutputCols) += 1;
953,063,715✔
2533
    } else if (info != NULL) {
9,232,059✔
2534
      // select distinct tbname from stb where tbname='abc';
2535
      info->needOutput = false;
9,146,366✔
2536
    }
2537
  }
2538

2539
  pMatchInfo->pList = pList;
269,220,889✔
2540

2541
_end:
269,188,336✔
2542
  if (code != TSDB_CODE_SUCCESS) {
269,188,336✔
2543
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
2544
  }
2545
  return code;
269,178,693✔
2546
}
2547

2548
static SResSchema createResSchema(int32_t type, int32_t bytes, int32_t slotId, int32_t scale, int32_t precision,
816,089,065✔
2549
                                  const char* name) {
2550
  SResSchema s = {0};
816,089,065✔
2551
  s.scale = scale;
816,137,498✔
2552
  s.type = type;
816,137,498✔
2553
  s.bytes = bytes;
816,137,498✔
2554
  s.slotId = slotId;
816,137,498✔
2555
  s.precision = precision;
816,137,498✔
2556
  tstrncpy(s.name, name, tListLen(s.name));
816,137,498✔
2557

2558
  return s;
816,137,498✔
2559
}
2560

2561
static SColumn* createColumn(int32_t blockId, int32_t slotId, int32_t colId, SDataType* pType, EColumnType colType) {
731,018,605✔
2562
  SColumn* pCol = taosMemoryCalloc(1, sizeof(SColumn));
731,018,605✔
2563
  if (pCol == NULL) {
730,702,589✔
2564
    return NULL;
×
2565
  }
2566

2567
  pCol->slotId = slotId;
730,702,589✔
2568
  pCol->colId = colId;
730,729,898✔
2569
  pCol->bytes = pType->bytes;
730,744,404✔
2570
  pCol->type = pType->type;
730,915,573✔
2571
  pCol->scale = pType->scale;
730,898,615✔
2572
  pCol->precision = pType->precision;
731,051,467✔
2573
  pCol->dataBlockId = blockId;
730,922,155✔
2574
  pCol->colType = colType;
730,880,800✔
2575
  return pCol;
731,014,754✔
2576
}
2577

2578
int32_t createExprFromOneNode(SExprInfo* pExp, SNode* pNode, int16_t slotId) {
820,822,796✔
2579
  int32_t code = TSDB_CODE_SUCCESS;
820,822,796✔
2580
  int32_t lino = 0;
820,822,796✔
2581
  pExp->base.numOfParams = 0;
820,822,796✔
2582
  pExp->base.pParam = NULL;
820,860,905✔
2583
  pExp->pExpr = taosMemoryCalloc(1, sizeof(tExprNode));
820,712,173✔
2584
  QUERY_CHECK_NULL(pExp->pExpr, code, lino, _end, terrno);
820,387,091✔
2585

2586
  pExp->pExpr->_function.num = 1;
820,498,405✔
2587
  pExp->pExpr->_function.functionId = -1;
820,631,759✔
2588

2589
  int32_t type = nodeType(pNode);
820,713,041✔
2590
  // it is a project query, or group by column
2591
  if (type == QUERY_NODE_COLUMN) {
820,760,382✔
2592
    pExp->pExpr->nodeType = QUERY_NODE_COLUMN;
497,393,604✔
2593
    SColumnNode* pColNode = (SColumnNode*)pNode;
497,418,498✔
2594

2595
    pExp->base.pParam = taosMemoryCalloc(1, sizeof(SFunctParam));
497,418,498✔
2596
    QUERY_CHECK_NULL(pExp->base.pParam, code, lino, _end, terrno);
497,320,573✔
2597

2598
    pExp->base.numOfParams = 1;
497,268,081✔
2599

2600
    SDataType* pType = &pColNode->node.resType;
497,403,074✔
2601
    pExp->base.resSchema =
2602
        createResSchema(pType->type, pType->bytes, slotId, pType->scale, pType->precision, pColNode->colName);
497,357,301✔
2603

2604
    pExp->base.pParam[0].pCol =
994,908,089✔
2605
        createColumn(pColNode->dataBlockId, pColNode->slotId, pColNode->colId, pType, pColNode->colType);
994,770,502✔
2606
    QUERY_CHECK_NULL(pExp->base.pParam[0].pCol, code, lino, _end, terrno);
497,464,340✔
2607

2608
    pExp->base.pParam[0].type = FUNC_PARAM_TYPE_COLUMN;
497,420,561✔
2609
  } else if (type == QUERY_NODE_VALUE) {
323,366,778✔
2610
    pExp->pExpr->nodeType = QUERY_NODE_VALUE;
16,224,907✔
2611
    SValueNode* pValNode = (SValueNode*)pNode;
16,225,783✔
2612

2613
    pExp->base.pParam = taosMemoryCalloc(1, sizeof(SFunctParam));
16,225,783✔
2614
    QUERY_CHECK_NULL(pExp->base.pParam, code, lino, _end, terrno);
16,221,143✔
2615

2616
    pExp->base.numOfParams = 1;
16,214,050✔
2617

2618
    SDataType* pType = &pValNode->node.resType;
16,224,602✔
2619
    pExp->base.resSchema =
2620
        createResSchema(pType->type, pType->bytes, slotId, pType->scale, pType->precision, pValNode->node.aliasName);
16,221,992✔
2621
    pExp->base.pParam[0].type = FUNC_PARAM_TYPE_VALUE;
16,217,406✔
2622
    code = nodesValueNodeToVariant(pValNode, &pExp->base.pParam[0].param);
16,222,824✔
2623
    QUERY_CHECK_CODE(code, lino, _end);
16,222,926✔
2624
  } else if (type == QUERY_NODE_REMOTE_VALUE) {
307,141,871✔
2625
    SRemoteValueNode* pRemote = (SRemoteValueNode*)pNode;
28,323,137✔
2626
    code = qFetchRemoteValue(gTaskScalarExtra.pSubJobCtx, pRemote->subQIdx, pRemote);
28,323,137✔
2627
    QUERY_CHECK_CODE(code, lino, _end);
28,334,378✔
2628

2629
    pExp->pExpr->nodeType = QUERY_NODE_VALUE;
23,582,754✔
2630
    SValueNode* pValNode = (SValueNode*)pNode;
23,582,215✔
2631

2632
    pExp->base.pParam = taosMemoryCalloc(1, sizeof(SFunctParam));
23,582,215✔
2633
    QUERY_CHECK_NULL(pExp->base.pParam, code, lino, _end, terrno);
23,582,767✔
2634

2635
    pExp->base.numOfParams = 1;
23,582,164✔
2636

2637
    SDataType* pType = &pValNode->node.resType;
23,582,213✔
2638
    pExp->base.resSchema =
2639
        createResSchema(pType->type, pType->bytes, slotId, pType->scale, pType->precision, pValNode->node.aliasName);
23,582,202✔
2640
    pExp->base.pParam[0].type = FUNC_PARAM_TYPE_VALUE;
23,582,718✔
2641
    code = nodesValueNodeToVariant(pValNode, &pExp->base.pParam[0].param);
23,581,612✔
2642
    QUERY_CHECK_CODE(code, lino, _end);
23,581,659✔
2643
  } else if (type == QUERY_NODE_FUNCTION) {
278,818,734✔
2644
    pExp->pExpr->nodeType = QUERY_NODE_FUNCTION;
247,217,260✔
2645
    SFunctionNode* pFuncNode = (SFunctionNode*)pNode;
247,221,583✔
2646

2647
    SDataType* pType = &pFuncNode->node.resType;
247,221,583✔
2648
    pExp->base.resSchema =
2649
        createResSchema(pType->type, pType->bytes, slotId, pType->scale, pType->precision, pFuncNode->node.aliasName);
247,226,844✔
2650
    tExprNode* pExprNode = pExp->pExpr;
247,190,307✔
2651

2652
    pExprNode->_function.functionId = pFuncNode->funcId;
247,200,831✔
2653
    pExprNode->_function.pFunctNode = pFuncNode;
247,223,622✔
2654
    pExprNode->_function.functionType = pFuncNode->funcType;
247,209,460✔
2655

2656
    tstrncpy(pExprNode->_function.functionName, pFuncNode->functionName, tListLen(pExprNode->_function.functionName));
247,240,803✔
2657

2658
    pExp->base.pParamList = pFuncNode->pParameterList;
247,248,688✔
2659
#if 1
2660
    // todo refactor: add the parameter for tbname function
2661
    const char* name = "tbname";
247,234,206✔
2662
    int32_t     len = strlen(name);
247,234,206✔
2663

2664
    if (!pFuncNode->pParameterList && (memcmp(pExprNode->_function.functionName, name, len) == 0) &&
247,234,206✔
2665
        pExprNode->_function.functionName[len] == 0) {
8,164,176✔
2666
      pFuncNode->pParameterList = NULL;
8,162,944✔
2667
      int32_t     code = nodesMakeList(&pFuncNode->pParameterList);
8,163,920✔
2668
      SValueNode* res = NULL;
8,164,921✔
2669
      if (TSDB_CODE_SUCCESS == code) {
8,161,779✔
2670
        code = nodesMakeNode(QUERY_NODE_VALUE, (SNode**)&res);
8,163,928✔
2671
      }
2672
      QUERY_CHECK_CODE(code, lino, _end);
8,166,302✔
2673
      res->node.resType = (SDataType){.bytes = sizeof(int64_t), .type = TSDB_DATA_TYPE_BIGINT};
8,166,302✔
2674
      code = nodesListAppend(pFuncNode->pParameterList, (SNode*)res);
8,166,406✔
2675
      if (code != TSDB_CODE_SUCCESS) {
8,167,313✔
2676
        nodesDestroyNode((SNode*)res);
×
2677
        res = NULL;
×
2678
      }
2679
      QUERY_CHECK_CODE(code, lino, _end);
8,167,313✔
2680
    }
2681
#endif
2682

2683
    int32_t numOfParam = LIST_LENGTH(pFuncNode->pParameterList);
247,256,575✔
2684

2685
    pExp->base.pParam = taosMemoryCalloc(numOfParam, sizeof(SFunctParam));
247,235,705✔
2686
    QUERY_CHECK_NULL(pExp->base.pParam, code, lino, _end, terrno);
247,178,355✔
2687
    pExp->base.numOfParams = numOfParam;
247,212,744✔
2688

2689
    for (int32_t j = 0; j < numOfParam && TSDB_CODE_SUCCESS == code; ++j) {
612,019,646✔
2690
      SNode* p1 = nodesListGetNode(pFuncNode->pParameterList, j);
365,021,155✔
2691
      QUERY_CHECK_NULL(p1, code, lino, _end, terrno);
365,040,628✔
2692
      if (p1->type == QUERY_NODE_COLUMN) {
365,040,628✔
2693
        SColumnNode* pcn = (SColumnNode*)p1;
233,583,452✔
2694

2695
        pExp->base.pParam[j].type = FUNC_PARAM_TYPE_COLUMN;
233,583,452✔
2696
        pExp->base.pParam[j].pCol =
467,145,667✔
2697
            createColumn(pcn->dataBlockId, pcn->slotId, pcn->colId, &pcn->node.resType, pcn->colType);
467,133,503✔
2698
        QUERY_CHECK_NULL(pExp->base.pParam[j].pCol, code, lino, _end, terrno);
233,580,360✔
2699
      } else if (p1->type == QUERY_NODE_VALUE) {
131,458,000✔
2700
        SValueNode* pvn = (SValueNode*)p1;
66,745,063✔
2701
        pExp->base.pParam[j].type = FUNC_PARAM_TYPE_VALUE;
66,745,063✔
2702
        code = nodesValueNodeToVariant(pvn, &pExp->base.pParam[j].param);
66,737,672✔
2703
        QUERY_CHECK_CODE(code, lino, _end);
66,743,441✔
2704
      } else if (p1->type == QUERY_NODE_REMOTE_VALUE) {
64,706,062✔
2705
        SRemoteValueNode* pRemote = (SRemoteValueNode*)p1;
1,435,022✔
2706
        code = qFetchRemoteValue(gTaskScalarExtra.pSubJobCtx, pRemote->subQIdx, pRemote);
1,435,022✔
2707
        QUERY_CHECK_CODE(code, lino, _end);
1,435,022✔
2708

2709
        SValueNode* pvn = (SValueNode*)pRemote;
1,218,033✔
2710
        pExp->base.pParam[j].type = FUNC_PARAM_TYPE_VALUE;
1,218,033✔
2711
        code = nodesValueNodeToVariant(pvn, &pExp->base.pParam[j].param);
1,218,033✔
2712
        QUERY_CHECK_CODE(code, lino, _end);
1,217,606✔
2713
      }
2714
    }
2715
    pExp->pExpr->_function.bindExprID = ((SExprNode*)pNode)->bindExprID;
246,998,491✔
2716
  } else if (type == QUERY_NODE_OPERATOR) {
31,601,474✔
2717
    pExp->pExpr->nodeType = QUERY_NODE_OPERATOR;
25,922,367✔
2718
    SOperatorNode* pOpNode = (SOperatorNode*)pNode;
25,922,510✔
2719

2720
    pExp->base.pParam = taosMemoryCalloc(1, sizeof(SFunctParam));
25,922,510✔
2721
    QUERY_CHECK_NULL(pExp->base.pParam, code, lino, _end, terrno);
25,921,450✔
2722
    pExp->base.numOfParams = 1;
25,921,607✔
2723

2724
    SDataType* pType = &pOpNode->node.resType;
25,924,435✔
2725
    pExp->base.resSchema =
2726
        createResSchema(pType->type, pType->bytes, slotId, pType->scale, pType->precision, pOpNode->node.aliasName);
25,921,749✔
2727
    pExp->pExpr->_optrRoot.pRootNode = pNode;
25,922,481✔
2728
  } else if (type == QUERY_NODE_CASE_WHEN) {
5,679,107✔
2729
    pExp->pExpr->nodeType = QUERY_NODE_OPERATOR;
5,677,398✔
2730
    SCaseWhenNode* pCaseNode = (SCaseWhenNode*)pNode;
5,677,398✔
2731

2732
    pExp->base.pParam = taosMemoryCalloc(1, sizeof(SFunctParam));
5,677,398✔
2733
    QUERY_CHECK_NULL(pExp->base.pParam, code, lino, _end, terrno);
5,677,398✔
2734
    pExp->base.numOfParams = 1;
5,677,398✔
2735

2736
    SDataType* pType = &pCaseNode->node.resType;
5,677,398✔
2737
    pExp->base.resSchema =
2738
        createResSchema(pType->type, pType->bytes, slotId, pType->scale, pType->precision, pCaseNode->node.aliasName);
5,677,398✔
2739
    pExp->pExpr->_optrRoot.pRootNode = pNode;
5,677,398✔
2740
  } else if (type == QUERY_NODE_LOGIC_CONDITION) {
2,128✔
2741
    pExp->pExpr->nodeType = QUERY_NODE_OPERATOR;
1,162✔
2742
    SLogicConditionNode* pCond = (SLogicConditionNode*)pNode;
1,162✔
2743
    pExp->base.pParam = taosMemoryCalloc(1, sizeof(SFunctParam));
1,162✔
2744
    QUERY_CHECK_NULL(pExp->base.pParam, code, lino, _end, terrno);
1,162✔
2745
    pExp->base.numOfParams = 1;
1,162✔
2746
    SDataType* pType = &pCond->node.resType;
1,162✔
2747
    pExp->base.resSchema =
2748
        createResSchema(pType->type, pType->bytes, slotId, pType->scale, pType->precision, pCond->node.aliasName);
1,162✔
2749
    pExp->pExpr->_optrRoot.pRootNode = pNode;
1,162✔
2750
  } else {
2751
    code = TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR;
1,144✔
2752
    QUERY_CHECK_CODE(code, lino, _end);
1,144✔
2753
  }
2754
  pExp->pExpr->relatedTo = ((SExprNode*)pNode)->relatedTo;
815,797,364✔
2755
_end:
820,792,997✔
2756
  if (code != TSDB_CODE_SUCCESS) {
820,792,997✔
2757
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
4,969,178✔
2758
  }
2759
  return code;
820,717,878✔
2760
}
2761

2762
int32_t createExprFromTargetNode(SExprInfo* pExp, STargetNode* pTargetNode) {
820,777,533✔
2763
  return createExprFromOneNode(pExp, pTargetNode->pExpr, pTargetNode->slotId);
820,777,533✔
2764
}
2765

2766
SExprInfo* createExpr(SNodeList* pNodeList, int32_t* numOfExprs) {
×
2767
  *numOfExprs = LIST_LENGTH(pNodeList);
×
2768
  SExprInfo* pExprs = taosMemoryCalloc(*numOfExprs, sizeof(SExprInfo));
×
2769
  if (!pExprs) {
×
2770
    return NULL;
×
2771
  }
2772

2773
  for (int32_t i = 0; i < (*numOfExprs); ++i) {
×
2774
    SExprInfo* pExp = &pExprs[i];
×
2775
    int32_t    code = createExprFromOneNode(pExp, nodesListGetNode(pNodeList, i), i + UD_TAG_COLUMN_INDEX);
×
2776
    if (code != TSDB_CODE_SUCCESS) {
×
2777
      taosMemoryFreeClear(pExprs);
×
2778
      terrno = code;
×
2779
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(code));
×
2780
      return NULL;
×
2781
    }
2782
  }
2783

2784
  return pExprs;
×
2785
}
2786

2787
int32_t createExprInfo(SNodeList* pNodeList, SNodeList* pGroupKeys, SExprInfo** pExprInfo, int32_t* numOfExprs) {
363,119,541✔
2788
  QRY_PARAM_CHECK(pExprInfo);
363,119,541✔
2789

2790
  int32_t code = 0;
363,166,779✔
2791
  int32_t numOfFuncs = LIST_LENGTH(pNodeList);
363,166,779✔
2792
  int32_t numOfGroupKeys = 0;
363,111,394✔
2793
  if (pGroupKeys != NULL) {
363,111,394✔
2794
    numOfGroupKeys = LIST_LENGTH(pGroupKeys);
35,215,017✔
2795
  }
2796

2797
  *numOfExprs = numOfFuncs + numOfGroupKeys;
363,113,602✔
2798
  if (*numOfExprs == 0) {
363,120,051✔
2799
    return code;
45,122,147✔
2800
  }
2801

2802
  SExprInfo* pExprs = taosMemoryCalloc(*numOfExprs, sizeof(SExprInfo));
318,014,912✔
2803
  if (pExprs == NULL) {
317,685,841✔
2804
    return terrno;
×
2805
  }
2806

2807
  for (int32_t i = 0; i < (*numOfExprs); ++i) {
1,133,254,440✔
2808
    STargetNode* pTargetNode = NULL;
820,549,239✔
2809
    if (i < numOfFuncs) {
820,549,239✔
2810
      pTargetNode = (STargetNode*)nodesListGetNode(pNodeList, i);
774,012,798✔
2811
    } else {
2812
      pTargetNode = (STargetNode*)nodesListGetNode(pGroupKeys, i - numOfFuncs);
46,536,441✔
2813
    }
2814
    if (!pTargetNode) {
820,670,360✔
2815
      destroyExprInfo(pExprs, *numOfExprs);
×
2816
      taosMemoryFreeClear(pExprs);
×
2817
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
2818
      return terrno;
×
2819
    }
2820

2821
    SExprInfo* pExp = &pExprs[i];
820,670,360✔
2822
    code = createExprFromTargetNode(pExp, pTargetNode);
820,700,124✔
2823
    if (code != TSDB_CODE_SUCCESS) {
820,537,777✔
2824
      destroyExprInfo(pExprs, *numOfExprs);
4,969,178✔
2825
      taosMemoryFreeClear(pExprs);
4,969,178✔
2826
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(code));
4,968,613✔
2827
      return code;
4,969,178✔
2828
    }
2829
  }
2830

2831
  *pExprInfo = pExprs;
312,945,056✔
2832
  return code;
313,036,224✔
2833
}
2834

2835
static void deleteSubsidiareCtx(void* pData) {
×
2836
  SSubsidiaryResInfo* pCtx = (SSubsidiaryResInfo*)pData;
×
2837
  if (pCtx->pCtx) {
×
2838
    taosMemoryFreeClear(pCtx->pCtx);
×
2839
  }
2840
}
×
2841

2842
// set the output buffer for the selectivity + tag query
2843
static int32_t setSelectValueColumnInfo(SqlFunctionCtx* pCtx, int32_t numOfOutput) {
335,753,222✔
2844
  int32_t num = 0;
335,753,222✔
2845
  int32_t code = TSDB_CODE_SUCCESS;
335,753,222✔
2846
  int32_t lino = 0;
335,753,222✔
2847

2848
  SArray* pValCtxArray = NULL;
335,753,222✔
2849
  for (int32_t i = numOfOutput - 1; i > 0; --i) {  // select Func is at the end of the list
833,599,647✔
2850
    int32_t funcIdx = pCtx[i].pExpr->pExpr->_function.bindExprID;
497,879,123✔
2851
    if (funcIdx > 0) {
497,861,236✔
2852
      if (pValCtxArray == NULL) {
37,266✔
2853
        // the end of the list is the select function of biggest index
2854
        pValCtxArray = taosArrayInit_s(sizeof(SSubsidiaryResInfo*), funcIdx);
24,866✔
2855
        if (pValCtxArray == NULL) {
24,866✔
2856
          return terrno;
×
2857
        }
2858
      }
2859
      if (funcIdx > pValCtxArray->size) {
37,266✔
2860
        qError("funcIdx:%d is out of range", funcIdx);
×
2861
        taosArrayDestroyP(pValCtxArray, deleteSubsidiareCtx);
×
2862
        return TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR;
×
2863
      }
2864
      SSubsidiaryResInfo* pSubsidiary = &pCtx[i].subsidiaries;
37,266✔
2865
      pSubsidiary->pCtx = taosMemoryCalloc(numOfOutput, POINTER_BYTES);
37,266✔
2866
      if (pSubsidiary->pCtx == NULL) {
37,266✔
2867
        taosArrayDestroyP(pValCtxArray, deleteSubsidiareCtx);
×
2868
        return terrno;
×
2869
      }
2870
      pSubsidiary->num = 0;
37,266✔
2871
      taosArraySet(pValCtxArray, funcIdx - 1, &pSubsidiary);
37,266✔
2872
    }
2873
  }
2874

2875
  SqlFunctionCtx*  p = NULL;
335,720,524✔
2876
  SqlFunctionCtx** pValCtx = NULL;
335,720,524✔
2877
  if (pValCtxArray == NULL) {
335,720,524✔
2878
    pValCtx = taosMemoryCalloc(numOfOutput, POINTER_BYTES);
335,720,028✔
2879
    if (pValCtx == NULL) {
335,570,826✔
2880
      QUERY_CHECK_CODE(terrno, lino, _end);
×
2881
    }
2882
  }
2883

2884
  for (int32_t i = 0; i < numOfOutput; ++i) {
1,136,968,320✔
2885
    const char* pName = pCtx[i].pExpr->pExpr->_function.functionName;
801,388,380✔
2886
    if ((strcmp(pName, "_select_value") == 0)) {
801,345,330✔
2887
      if (pValCtxArray == NULL) {
3,276,285✔
2888
        pValCtx[num++] = &pCtx[i];
3,229,718✔
2889
      } else {
2890
        int32_t bindFuncIndex = pCtx[i].pExpr->pExpr->relatedTo;  // start from index 1;
46,567✔
2891
        if (bindFuncIndex > 0) {                                  // 0 is default index related to the select function
46,567✔
2892
          bindFuncIndex -= 1;
46,567✔
2893
        }
2894
        SSubsidiaryResInfo** pSubsidiary = taosArrayGet(pValCtxArray, bindFuncIndex);
46,567✔
2895
        if (pSubsidiary == NULL) {
46,567✔
2896
          QUERY_CHECK_CODE(TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR, lino, _end);
×
2897
        }
2898
        (*pSubsidiary)->pCtx[(*pSubsidiary)->num] = &pCtx[i];
46,567✔
2899
        (*pSubsidiary)->num++;
46,567✔
2900
      }
2901
    } else if (fmIsSelectFunc(pCtx[i].functionId)) {
798,069,045✔
2902
      if (pValCtxArray == NULL) {
54,400,251✔
2903
        p = &pCtx[i];
54,348,123✔
2904
      }
2905
    }
2906
  }
2907

2908
  if (p != NULL) {
335,579,940✔
2909
    p->subsidiaries.pCtx = pValCtx;
24,607,775✔
2910
    p->subsidiaries.num = num;
24,603,794✔
2911
  } else {
2912
    taosMemoryFreeClear(pValCtx);
310,972,165✔
2913
  }
2914

UNCOV
2915
_end:
×
2916
  if (code != TSDB_CODE_SUCCESS) {
335,595,435✔
2917
    taosArrayDestroyP(pValCtxArray, deleteSubsidiareCtx);
×
2918
    taosMemoryFreeClear(pValCtx);
×
2919
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
2920
  } else {
2921
    taosArrayDestroy(pValCtxArray);
335,595,435✔
2922
  }
2923
  return code;
335,662,798✔
2924
}
2925

2926
SqlFunctionCtx* createSqlFunctionCtx(SExprInfo* pExprInfo, int32_t numOfOutput, int32_t** rowEntryInfoOffset,
335,775,052✔
2927
                                     SFunctionStateStore* pStore) {
2928
  int32_t         code = TSDB_CODE_SUCCESS;
335,775,052✔
2929
  int32_t         lino = 0;
335,775,052✔
2930
  SqlFunctionCtx* pFuncCtx = (SqlFunctionCtx*)taosMemoryCalloc(numOfOutput, sizeof(SqlFunctionCtx));
335,775,052✔
2931
  if (pFuncCtx == NULL) {
335,455,688✔
2932
    return NULL;
×
2933
  }
2934

2935
  *rowEntryInfoOffset = taosMemoryCalloc(numOfOutput, sizeof(int32_t));
335,455,688✔
2936
  if (*rowEntryInfoOffset == 0) {
335,735,524✔
2937
    taosMemoryFreeClear(pFuncCtx);
×
2938
    return NULL;
×
2939
  }
2940

2941
  for (int32_t i = 0; i < numOfOutput; ++i) {
1,137,184,826✔
2942
    SExprInfo* pExpr = &pExprInfo[i];
801,540,240✔
2943

2944
    SExprBasicInfo* pFunct = &pExpr->base;
801,440,629✔
2945
    SqlFunctionCtx* pCtx = &pFuncCtx[i];
801,490,826✔
2946

2947
    pCtx->functionId = -1;
801,499,798✔
2948
    pCtx->pExpr = pExpr;
801,531,849✔
2949

2950
    if (pExpr->pExpr->nodeType == QUERY_NODE_FUNCTION) {
801,510,302✔
2951
      SFuncExecEnv env = {0};
245,674,652✔
2952
      pCtx->functionId = pExpr->pExpr->_function.pFunctNode->funcId;
245,673,719✔
2953
      pCtx->isPseudoFunc = fmIsWindowPseudoColumnFunc(pCtx->functionId) || fmIsPlaceHolderFunc(pCtx->functionId);
245,674,910✔
2954
      pCtx->isNotNullFunc = fmIsNotNullOutputFunc(pCtx->functionId);
245,675,566✔
2955

2956
      bool isUdaf = fmIsUserDefinedFunc(pCtx->functionId);
245,690,723✔
2957
      if (fmIsAggFunc(pCtx->functionId) || fmIsIndefiniteRowsFunc(pCtx->functionId)) {
403,875,727✔
2958
        if (!isUdaf) {
158,210,281✔
2959
          code = fmGetFuncExecFuncs(pCtx->functionId, &pCtx->fpSet);
158,168,044✔
2960
          QUERY_CHECK_CODE(code, lino, _end);
158,163,349✔
2961
        } else {
2962
          char* udfName = pExpr->pExpr->_function.pFunctNode->functionName;
42,237✔
2963
          pCtx->udfName = taosStrdup(udfName);
42,237✔
2964
          QUERY_CHECK_NULL(pCtx->udfName, code, lino, _end, terrno);
42,237✔
2965

2966
          code = fmGetUdafExecFuncs(pCtx->functionId, &pCtx->fpSet);
42,237✔
2967
          QUERY_CHECK_CODE(code, lino, _end);
42,237✔
2968
        }
2969
        bool tmp = pCtx->fpSet.getEnv(pExpr->pExpr->_function.pFunctNode, &env);
158,205,586✔
2970
        if (!tmp) {
158,203,967✔
2971
          code = terrno;
×
2972
          QUERY_CHECK_CODE(code, lino, _end);
59✔
2973
        }
2974
      } else {
2975
        if (fmIsPlaceHolderFunc(pCtx->functionId)) {
87,451,230✔
2976
          code = fmGetStreamPesudoFuncEnv(pCtx->functionId, pExpr->base.pParamList, &env);
10,123,234✔
2977
          QUERY_CHECK_CODE(code, lino, _end);
10,124,483✔
2978
        }      
2979
        
2980
        code = fmGetScalarFuncExecFuncs(pCtx->functionId, &pCtx->sfp);
87,457,733✔
2981
        if (code != TSDB_CODE_SUCCESS && isUdaf) {
87,455,441✔
2982
          code = TSDB_CODE_SUCCESS;
26,100✔
2983
        }
2984
        QUERY_CHECK_CODE(code, lino, _end);
87,455,441✔
2985

2986
        if (pCtx->sfp.getEnv != NULL) {
87,455,441✔
2987
          bool tmp = pCtx->sfp.getEnv(pExpr->pExpr->_function.pFunctNode, &env);
18,425,520✔
2988
          if (!tmp) {
18,423,188✔
2989
            code = terrno;
×
2990
            QUERY_CHECK_CODE(code, lino, _end);
×
2991
          }
2992
        }
2993
      }
2994
      pCtx->resDataInfo.interBufSize = env.calcMemSize;
245,656,972✔
2995
    } else if (pExpr->pExpr->nodeType == QUERY_NODE_COLUMN || pExpr->pExpr->nodeType == QUERY_NODE_OPERATOR ||
555,720,790✔
2996
               pExpr->pExpr->nodeType == QUERY_NODE_VALUE) {
39,755,273✔
2997
      // for simple column, the result buffer needs to hold at least one element.
2998
      pCtx->resDataInfo.interBufSize = pFunct->resSchema.bytes;
555,887,738✔
2999
    }
3000

3001
    pCtx->input.numOfInputCols = pFunct->numOfParams;
801,556,317✔
3002
    pCtx->input.pData = taosMemoryCalloc(pFunct->numOfParams, POINTER_BYTES);
801,576,057✔
3003
    QUERY_CHECK_NULL(pCtx->input.pData, code, lino, _end, terrno);
801,450,376✔
3004
    pCtx->input.pColumnDataAgg = taosMemoryCalloc(pFunct->numOfParams, POINTER_BYTES);
801,540,231✔
3005
    QUERY_CHECK_NULL(pCtx->input.pColumnDataAgg, code, lino, _end, terrno);
801,449,474✔
3006

3007
    pCtx->pTsOutput = NULL;
801,538,935✔
3008
    pCtx->resDataInfo.bytes = pFunct->resSchema.bytes;
801,534,597✔
3009
    pCtx->resDataInfo.type = pFunct->resSchema.type;
801,594,379✔
3010
    pCtx->order = TSDB_ORDER_ASC;
801,543,857✔
3011
    pCtx->start.key = INT64_MIN;
801,432,579✔
3012
    pCtx->end.key = INT64_MIN;
801,570,657✔
3013
    pCtx->numOfParams = pExpr->base.numOfParams;
801,613,709✔
3014
    pCtx->param = pFunct->pParam;
801,501,255✔
3015
    pCtx->saveHandle.currentPage = -1;
801,660,286✔
3016
    pCtx->pStore = pStore;
801,613,571✔
3017
    pCtx->hasWindowOrGroup = false;
801,515,602✔
3018
    pCtx->needCleanup = false;
801,347,355✔
3019
    pCtx->skipDynDataCheck = false;
801,489,858✔
3020
  }
3021

3022
  for (int32_t i = 1; i < numOfOutput; ++i) {
833,666,418✔
3023
    (*rowEntryInfoOffset)[i] = (int32_t)((*rowEntryInfoOffset)[i - 1] + sizeof(SResultRowEntryInfo) +
995,795,695✔
3024
                                         pFuncCtx[i - 1].resDataInfo.interBufSize);
497,909,440✔
3025
  }
3026

3027
  code = setSelectValueColumnInfo(pFuncCtx, numOfOutput);
335,761,949✔
3028
  QUERY_CHECK_CODE(code, lino, _end);
335,662,685✔
3029

3030
_end:
335,662,685✔
3031
  if (code != TSDB_CODE_SUCCESS) {
335,652,208✔
3032
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
3033
    for (int32_t i = 0; i < numOfOutput; ++i) {
×
3034
      taosMemoryFree(pFuncCtx[i].input.pData);
×
UNCOV
3035
      taosMemoryFree(pFuncCtx[i].input.pColumnDataAgg);
×
3036
    }
3037
    taosMemoryFreeClear(*rowEntryInfoOffset);
×
UNCOV
3038
    taosMemoryFreeClear(pFuncCtx);
×
3039

3040
    terrno = code;
×
UNCOV
3041
    return NULL;
×
3042
  }
3043
  return pFuncCtx;
335,652,208✔
3044
}
3045

3046
// NOTE: sources columns are more than the destination SSDatablock columns.
3047
// doFilter in table scan needs every column even its output is false
3048
int32_t relocateColumnData(SSDataBlock* pBlock, const SArray* pColMatchInfo, SArray* pCols, bool outputEveryColumn) {
9,243,002✔
3049
  int32_t code = TSDB_CODE_SUCCESS;
9,243,002✔
3050
  size_t  numOfSrcCols = taosArrayGetSize(pCols);
9,243,002✔
3051

3052
  int32_t i = 0, j = 0;
9,243,703✔
3053
  while (i < numOfSrcCols && j < taosArrayGetSize(pColMatchInfo)) {
85,253,939✔
3054
    SColumnInfoData* p = taosArrayGet(pCols, i);
76,010,981✔
3055
    if (!p) {
76,010,981✔
3056
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
UNCOV
3057
      return terrno;
×
3058
    }
3059
    SColMatchItem* pmInfo = taosArrayGet(pColMatchInfo, j);
76,010,981✔
3060
    if (!pmInfo) {
76,010,937✔
UNCOV
3061
      return terrno;
×
3062
    }
3063

3064
    if (p->info.colId == pmInfo->colId) {
76,010,937✔
3065
      SColumnInfoData* pDst = taosArrayGet(pBlock->pDataBlock, pmInfo->dstSlotId);
67,413,074✔
3066
      if (!pDst) {
67,413,074✔
UNCOV
3067
        return terrno;
×
3068
      }
3069
      code = colDataAssign(pDst, p, pBlock->info.rows, &pBlock->info);
67,413,074✔
3070
      if (code != TSDB_CODE_SUCCESS) {
67,411,672✔
3071
        qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(code));
×
UNCOV
3072
        return code;
×
3073
      }
3074
      i++;
67,411,672✔
3075
      j++;
67,411,672✔
3076
    } else if (p->info.colId < pmInfo->colId) {
8,598,564✔
3077
      i++;
8,598,564✔
3078
    } else {
3079
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR));
×
UNCOV
3080
      return TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR;
×
3081
    }
3082
  }
3083
  return code;
9,243,703✔
3084
}
3085

3086
SInterval extractIntervalInfo(const STableScanPhysiNode* pTableScanNode) {
193,987,142✔
3087
  SInterval interval = {
387,963,888✔
3088
      .interval = pTableScanNode->interval,
193,898,097✔
3089
      .sliding = pTableScanNode->sliding,
193,908,766✔
3090
      .intervalUnit = pTableScanNode->intervalUnit,
194,025,842✔
3091
      .slidingUnit = pTableScanNode->slidingUnit,
193,950,874✔
3092
      .offset = pTableScanNode->offset,
193,796,041✔
3093
      .precision = pTableScanNode->scan.node.pOutputDataBlockDesc->precision,
193,827,145✔
3094
      .timeRange = pTableScanNode->scanRange,
3095
  };
3096
  calcIntervalAutoOffset(&interval);
194,032,893✔
3097

3098
  return interval;
193,740,729✔
3099
}
3100

3101
SColumn extractColumnFromColumnNode(SColumnNode* pColNode) {
51,612,693✔
3102
  SColumn c = {0};
51,612,693✔
3103

3104
  c.slotId = pColNode->slotId;
51,612,693✔
3105
  c.colId = pColNode->colId;
51,612,033✔
3106
  c.type = pColNode->node.resType.type;
51,612,007✔
3107
  c.bytes = pColNode->node.resType.bytes;
51,614,952✔
3108
  c.scale = pColNode->node.resType.scale;
51,609,527✔
3109
  c.precision = pColNode->node.resType.precision;
51,606,791✔
3110
  return c;
51,613,107✔
3111
}
3112

3113

3114
/**
3115
 * @brief Determine the actual time range for reading data based on the RANGE clause and the WHERE conditions.
3116
 * @param[in] cond The range specified by WHERE condition.
3117
 * @param[in] range The range specified by RANGE clause.
3118
 * @param[out] twindow The range to be read in DESC order, and only one record is needed.
3119
 * @param[out] extTwindow The external range to read for only one record, which is used for FILL clause.
3120
 * @note `cond` and `twindow` may be the same address.
3121
 */
3122
static int32_t getQueryExtWindow(const STimeWindow* cond, const STimeWindow* range, STimeWindow* twindow,
1,642,604✔
3123
                                 STimeWindow* extTwindows) {
3124
  int32_t     code = TSDB_CODE_SUCCESS;
1,642,604✔
3125
  int32_t     lino = 0;
1,642,604✔
3126
  STimeWindow tempWindow;
3127

3128
  if (cond->skey > cond->ekey || range->skey > range->ekey) {
1,642,604✔
3129
    *twindow = extTwindows[0] = extTwindows[1] = TSWINDOW_DESC_INITIALIZER;
2,980✔
3130
    return code;
2,980✔
3131
  }
3132

3133
  if (range->ekey < cond->skey) {
1,639,624✔
3134
    extTwindows[1] = *cond;
242,330✔
3135
    *twindow = extTwindows[0] = TSWINDOW_DESC_INITIALIZER;
242,330✔
3136
    return code;
242,330✔
3137
  }
3138

3139
  if (cond->ekey < range->skey) {
1,397,294✔
3140
    extTwindows[0] = *cond;
178,766✔
3141
    *twindow = extTwindows[1] = TSWINDOW_DESC_INITIALIZER;
178,766✔
3142
    return code;
178,766✔
3143
  }
3144

3145
  // Only scan data in the time range intersecion.
3146
  extTwindows[0] = extTwindows[1] = *cond;
1,218,528✔
3147
  twindow->skey = TMAX(cond->skey, range->skey);
1,218,528✔
3148
  twindow->ekey = TMIN(cond->ekey, range->ekey);
1,218,528✔
3149
  extTwindows[0].ekey = twindow->skey - 1;
1,218,528✔
3150
  extTwindows[1].skey = twindow->ekey + 1;
1,218,528✔
3151

3152
  return code;
1,218,528✔
3153
}
3154

3155
int32_t initQueryTableDataCond(SQueryTableDataCond* pCond, const STableScanPhysiNode* pTableScanNode,
217,263,248✔
3156
                               const SReadHandle* readHandle, bool applyExtWin) {
3157
  int32_t code = 0;                             
217,263,248✔
3158
  pCond->order = pTableScanNode->scanSeq[0] > 0 ? TSDB_ORDER_ASC : TSDB_ORDER_DESC;
217,263,248✔
3159
  pCond->numOfCols = LIST_LENGTH(pTableScanNode->scan.pScanCols);
217,283,572✔
3160

3161
  pCond->colList = taosMemoryCalloc(pCond->numOfCols, sizeof(SColumnInfo));
217,288,526✔
3162
  if (!pCond->colList) {
217,196,870✔
UNCOV
3163
    return terrno;
×
3164
  }
3165
  pCond->pSlotList = taosMemoryMalloc(sizeof(int32_t) * pCond->numOfCols);
217,165,724✔
3166
  if (pCond->pSlotList == NULL) {
217,220,845✔
3167
    taosMemoryFreeClear(pCond->colList);
×
UNCOV
3168
    return terrno;
×
3169
  }
3170

3171
  // TODO: get it from stable scan node
3172
  pCond->twindows = pTableScanNode->scanRange;
217,130,396✔
3173
  pCond->suid = pTableScanNode->scan.suid;
217,286,153✔
3174
  pCond->type = TIMEWINDOW_RANGE_CONTAINED;
217,123,029✔
3175
  pCond->startVersion = -1;
217,251,141✔
3176
  pCond->endVersion = -1;
217,293,962✔
3177
  pCond->skipRollup = readHandle->skipRollup;
217,073,474✔
3178
  if (readHandle->winRangeValid) {
217,215,186✔
3179
    pCond->twindows = readHandle->winRange;
300,560✔
3180
  }
3181
  pCond->cacheSttStatis = readHandle->cacheSttStatis;
217,297,706✔
3182
  // allowed read stt file optimization mode
3183
  pCond->notLoadData = (pTableScanNode->dataRequired == FUNC_DATA_REQUIRED_NOT_LOAD) &&
434,594,439✔
3184
                       (pTableScanNode->scan.node.pConditions == NULL) && (pTableScanNode->interval == 0);
217,247,769✔
3185

3186
  int32_t j = 0;
217,264,624✔
3187
  for (int32_t i = 0; i < pCond->numOfCols; ++i) {
1,027,374,902✔
3188
    STargetNode* pNode = (STargetNode*)nodesListGetNode(pTableScanNode->scan.pScanCols, i);
810,249,410✔
3189
    if (!pNode) {
809,884,420✔
3190
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
UNCOV
3191
      return terrno;
×
3192
    }
3193
    SColumnNode* pColNode = (SColumnNode*)pNode->pExpr;
809,884,420✔
3194
    if (pColNode->colType == COLUMN_TYPE_TAG) {
810,173,345✔
UNCOV
3195
      continue;
×
3196
    }
3197

3198
    pCond->colList[j].type = pColNode->node.resType.type;
810,131,787✔
3199
    pCond->colList[j].bytes = pColNode->node.resType.bytes;
810,128,227✔
3200
    pCond->colList[j].colId = pColNode->colId;
810,179,833✔
3201
    pCond->colList[j].pk = pColNode->isPk;
810,166,037✔
3202

3203
    pCond->pSlotList[j] = pNode->slotId;
810,226,323✔
3204
    j += 1;
810,110,278✔
3205
  }
3206

3207
  pCond->numOfCols = j;
217,333,460✔
3208

3209
  if (applyExtWin) {
217,337,124✔
3210
    if (NULL != pTableScanNode->pExtScanRange) {
194,367,097✔
3211
      pCond->type = TIMEWINDOW_RANGE_EXTERNAL;
1,642,604✔
3212
      code = getQueryExtWindow(&pCond->twindows, pTableScanNode->pExtScanRange, &pCond->twindows, pCond->extTwindows);
1,642,604✔
3213
    } else if (readHandle->extWinRangeValid) {
192,607,443✔
3214
      pCond->type = TIMEWINDOW_RANGE_EXTERNAL;
×
UNCOV
3215
      code = getQueryExtWindow(&pCond->twindows, &readHandle->extWinRange, &pCond->twindows, pCond->extTwindows);
×
3216
    }
3217
  }
3218
  
3219
  return code;
217,250,016✔
3220
}
3221

3222
int32_t initQueryTableDataCondWithColArray(SQueryTableDataCond* pCond, SQueryTableDataCond* pOrgCond,
3,062,378✔
3223
                                           const SReadHandle* readHandle, SArray* colArray) {
3224
  int32_t code = TSDB_CODE_SUCCESS;
3,062,378✔
3225
  int32_t lino = 0;
3,062,378✔
3226

3227
  pCond->order = TSDB_ORDER_ASC;
3,062,378✔
3228
  pCond->numOfCols = (int32_t)taosArrayGetSize(colArray);
3,062,378✔
3229

3230
  pCond->colList = taosMemoryCalloc(pCond->numOfCols, sizeof(SColumnInfo));
3,062,378✔
3231
  QUERY_CHECK_NULL(pCond->colList, code, lino, _return, terrno);
3,062,378✔
3232

3233
  pCond->pSlotList = taosMemoryMalloc(sizeof(int32_t) * pCond->numOfCols);
3,062,378✔
3234
  QUERY_CHECK_NULL(pCond->pSlotList, code, lino, _return, terrno);
3,062,378✔
3235

3236
  pCond->twindows = pOrgCond->twindows;
3,062,378✔
3237
  pCond->order = pOrgCond->order;
3,062,378✔
3238
  pCond->type = pOrgCond->type;
3,062,378✔
3239
  pCond->startVersion = -1;
3,062,378✔
3240
  pCond->endVersion = -1;
3,062,378✔
3241
  pCond->skipRollup = true;
3,062,378✔
3242
  pCond->notLoadData = false;
3,062,378✔
3243

3244
  for (int32_t i = 0; i < pCond->numOfCols; ++i) {
13,468,466✔
3245
    SColIdPair* pColPair = taosArrayGet(colArray, i);
10,406,088✔
3246
    QUERY_CHECK_NULL(pColPair, code, lino, _return, terrno);
10,406,088✔
3247

3248
    bool find = false;
10,406,088✔
3249
    for (int32_t j = 0; j < pOrgCond->numOfCols; ++j) {
56,406,474✔
3250
      if (pOrgCond->colList[j].colId == pColPair->vtbColId) {
56,406,474✔
3251
        pCond->colList[i].type = pOrgCond->colList[j].type;
10,406,088✔
3252
        pCond->colList[i].bytes = pOrgCond->colList[j].bytes;
10,406,088✔
3253
        pCond->colList[i].colId = pColPair->orgColId;
10,406,088✔
3254
        pCond->colList[i].pk = pOrgCond->colList[j].pk;
10,406,088✔
3255
        pCond->pSlotList[i] = i;
10,406,088✔
3256
        find = true;
10,406,088✔
3257
        break;
10,406,088✔
3258
      }
3259
    }
3260
    QUERY_CHECK_CONDITION(find, code, lino, _return, TSDB_CODE_NOT_FOUND);
10,406,088✔
3261
  }
3262

3263
  return code;
3,062,378✔
3264
_return:
×
3265
  qError("%s failed at line %d since %s", __func__, lino, tstrerror(terrno));
×
3266
  taosMemoryFreeClear(pCond->colList);
×
UNCOV
3267
  taosMemoryFreeClear(pCond->pSlotList);
×
UNCOV
3268
  return code;
×
3269
}
3270

3271
void cleanupQueryTableDataCond(SQueryTableDataCond* pCond) {
460,044,766✔
3272
  taosMemoryFreeClear(pCond->colList);
460,044,766✔
3273
  taosMemoryFreeClear(pCond->pSlotList);
460,016,715✔
3274
}
459,963,727✔
3275

3276
int32_t convertFillType(int32_t mode) {
2,067,372✔
3277
  int32_t type = TSDB_FILL_NONE;
2,067,372✔
3278
  switch (mode) {
2,067,372✔
3279
    case FILL_MODE_PREV:
111,455✔
3280
      type = TSDB_FILL_PREV;
111,455✔
3281
      break;
111,455✔
3282
    case FILL_MODE_NONE:
×
UNCOV
3283
      type = TSDB_FILL_NONE;
×
UNCOV
3284
      break;
×
3285
    case FILL_MODE_NULL:
124,906✔
3286
      type = TSDB_FILL_NULL;
124,906✔
3287
      break;
124,906✔
3288
    case FILL_MODE_NULL_F:
15,969✔
3289
      type = TSDB_FILL_NULL_F;
15,969✔
3290
      break;
15,969✔
3291
    case FILL_MODE_NEXT:
96,306✔
3292
      type = TSDB_FILL_NEXT;
96,306✔
3293
      break;
96,306✔
3294
    case FILL_MODE_VALUE:
149,105✔
3295
      type = TSDB_FILL_SET_VALUE;
149,105✔
3296
      break;
149,105✔
3297
    case FILL_MODE_VALUE_F:
4,290✔
3298
      type = TSDB_FILL_SET_VALUE_F;
4,290✔
3299
      break;
4,290✔
3300
    case FILL_MODE_LINEAR:
149,902✔
3301
      type = TSDB_FILL_LINEAR;
149,902✔
3302
      break;
149,902✔
3303
    case FILL_MODE_NEAR:
1,415,439✔
3304
      type = TSDB_FILL_NEAR;
1,415,439✔
3305
      break;
1,415,439✔
UNCOV
3306
    default:
×
UNCOV
3307
      type = TSDB_FILL_NONE;
×
3308
  }
3309

3310
  return type;
2,067,372✔
3311
}
3312

3313
void getInitialStartTimeWindow(SInterval* pInterval, TSKEY ts, STimeWindow* w, bool ascQuery) {
1,910,985,536✔
3314
  if (ascQuery) {
1,910,985,536✔
3315
    *w = getAlignQueryTimeWindow(pInterval, ts);
1,911,219,362✔
3316
  } else {
3317
    // the start position of the first time window in the endpoint that spreads beyond the queried last timestamp
3318
    *w = getAlignQueryTimeWindow(pInterval, ts);
56,181✔
3319

3320
    int64_t key = w->skey;
131,985✔
3321
    while (key < ts) {  // moving towards end
146,658✔
3322
      key = getNextTimeWindowStart(pInterval, key, TSDB_ORDER_ASC);
75,011✔
3323
      if (key > ts) {
75,011✔
3324
        break;
60,338✔
3325
      }
3326

3327
      w->skey = key;
14,673✔
3328
    }
3329
    w->ekey = taosTimeAdd(w->skey, pInterval->interval, pInterval->intervalUnit, pInterval->precision, NULL) - 1;
131,985✔
3330
  }
3331
}
1,910,565,265✔
3332

3333
static STimeWindow doCalculateTimeWindow(int64_t ts, SInterval* pInterval) {
27,522,136✔
3334
  STimeWindow w = {0};
27,522,136✔
3335

3336
  w.skey = taosTimeTruncate(ts, pInterval);
27,522,136✔
3337
  w.ekey = taosTimeGetIntervalEnd(w.skey, pInterval);
27,521,602✔
3338
  return w;
27,523,738✔
3339
}
3340

3341
STimeWindow getFirstQualifiedTimeWindow(int64_t ts, STimeWindow* pWindow, SInterval* pInterval, int32_t order) {
1,647,188✔
3342
  STimeWindow win = *pWindow;
1,647,188✔
3343
  STimeWindow save = win;
1,647,188✔
3344
  while (win.skey <= ts && win.ekey >= ts) {
9,191,257✔
3345
    save = win;
7,544,069✔
3346
    // get previous time window
3347
    getNextTimeWindow(pInterval, &win, order == TSDB_ORDER_DESC ? TSDB_ORDER_ASC : TSDB_ORDER_DESC);
7,544,069✔
3348
  }
3349

3350
  return save;
1,647,188✔
3351
}
3352

3353
// get the correct time window according to the handled timestamp
3354
// todo refactor
3355
STimeWindow getActiveTimeWindow(SDiskbasedBuf* pBuf, SResultRowInfo* pResultRowInfo, int64_t ts, SInterval* pInterval,
47,231,441✔
3356
                                int32_t order) {
3357
  STimeWindow w = {0};
47,231,441✔
3358
  if (pResultRowInfo->cur.pageId == -1) {  // the first window, from the previous stored value
47,231,759✔
3359
    getInitialStartTimeWindow(pInterval, ts, &w, (order != TSDB_ORDER_DESC));
6,005,597✔
3360
    return w;
6,008,305✔
3361
  }
3362

3363
  SResultRow* pRow = getResultRowByPos(pBuf, &pResultRowInfo->cur, false);
41,223,136✔
3364
  if (pRow) {
41,226,469✔
3365
    TAOS_SET_OBJ_ALIGNED(&w, pRow->win);
41,227,321✔
3366
  }
3367

3368
  // in case of typical time window, we can calculate time window directly.
3369
  if (w.skey > ts || w.ekey < ts) {
41,227,166✔
3370
    w = doCalculateTimeWindow(ts, pInterval);
27,522,769✔
3371
  }
3372

3373
  if (pInterval->interval != pInterval->sliding) {
41,228,135✔
3374
    // it is an sliding window query, in which sliding value is not equalled to
3375
    // interval value, and we need to find the first qualified time window.
3376
    w = getFirstQualifiedTimeWindow(ts, &w, pInterval, order);
1,647,188✔
3377
  }
3378

3379
  return w;
41,223,741✔
3380
}
3381

3382
TSKEY getNextTimeWindowStart(const SInterval* pInterval, TSKEY start, int32_t order) {
2,147,483,647✔
3383
  int32_t factor = GET_FORWARD_DIRECTION_FACTOR(order);
2,147,483,647✔
3384
  TSKEY   nextStart = taosTimeAdd(start, -1 * pInterval->offset, pInterval->offsetUnit, pInterval->precision, NULL);
2,147,483,647✔
3385
  nextStart = taosTimeAdd(nextStart, factor * pInterval->sliding, pInterval->slidingUnit, pInterval->precision, NULL);
2,147,483,647✔
3386
  nextStart = taosTimeAdd(nextStart, pInterval->offset, pInterval->offsetUnit, pInterval->precision, NULL);
2,147,483,647✔
3387
  return nextStart;
2,147,483,647✔
3388
}
3389

3390
void getNextTimeWindow(const SInterval* pInterval, STimeWindow* tw, int32_t order) {
2,147,483,647✔
3391
  tw->skey = getNextTimeWindowStart(pInterval, tw->skey, order);
2,147,483,647✔
3392
  tw->ekey = taosTimeAdd(tw->skey, pInterval->interval, pInterval->intervalUnit, pInterval->precision, NULL) - 1;
2,147,483,647✔
3393
}
2,147,483,647✔
3394

3395
bool hasLimitOffsetInfo(SLimitInfo* pLimitInfo) {
319,063,850✔
3396
  return (pLimitInfo->limit.limit != -1 || pLimitInfo->limit.offset != -1 || pLimitInfo->slimit.limit != -1 ||
635,638,351✔
3397
          pLimitInfo->slimit.offset != -1);
316,573,896✔
3398
}
3399

UNCOV
3400
bool hasSlimitOffsetInfo(SLimitInfo* pLimitInfo) {
×
UNCOV
3401
  return (pLimitInfo->slimit.limit != -1 || pLimitInfo->slimit.offset != -1);
×
3402
}
3403

3404
void initLimitInfo(const SNode* pLimit, const SNode* pSLimit, SLimitInfo* pLimitInfo) {
480,486,114✔
3405
  SLimit limit = {.limit = getLimit(pLimit), .offset = getOffset(pLimit)};
480,486,114✔
3406
  SLimit slimit = {.limit = getLimit(pSLimit), .offset = getOffset(pSLimit)};
480,282,170✔
3407

3408
  pLimitInfo->limit = limit;
480,325,736✔
3409
  pLimitInfo->slimit = slimit;
480,326,217✔
3410
  pLimitInfo->remainOffset = limit.offset;
480,312,268✔
3411
  pLimitInfo->remainGroupOffset = slimit.offset;
480,312,290✔
3412
  pLimitInfo->numOfOutputRows = 0;
480,311,474✔
3413
  pLimitInfo->numOfOutputGroups = 0;
480,374,854✔
3414
  pLimitInfo->currentGroupId = 0;
480,373,452✔
3415
}
480,465,773✔
3416

3417
void resetLimitInfoForNextGroup(SLimitInfo* pLimitInfo) {
52,057,493✔
3418
  pLimitInfo->numOfOutputRows = 0;
52,057,493✔
3419
  pLimitInfo->remainOffset = pLimitInfo->limit.offset;
52,063,867✔
3420
}
52,060,429✔
3421

3422
int32_t tableListGetSize(const STableListInfo* pTableList, int32_t* pRes) {
479,167,958✔
3423
  if (taosArrayGetSize(pTableList->pTableList) != taosHashGetSize(pTableList->map)) {
479,167,958✔
UNCOV
3424
    qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR));
×
UNCOV
3425
    return TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR;
×
3426
  }
3427
  (*pRes) = taosArrayGetSize(pTableList->pTableList);
479,177,594✔
3428
  return TSDB_CODE_SUCCESS;
479,189,838✔
3429
}
3430

3431
uint64_t tableListGetSuid(const STableListInfo* pTableList) { return pTableList->idInfo.suid; }
3,022,251✔
3432

3433
STableKeyInfo* tableListGetInfo(const STableListInfo* pTableList, int32_t index) {
155,718,645✔
3434
  if (taosArrayGetSize(pTableList->pTableList) == 0) {
155,718,645✔
3435
    return NULL;
3,677✔
3436
  }
3437

3438
  return taosArrayGet(pTableList->pTableList, index);
155,700,725✔
3439
}
3440

3441
int32_t tableListFind(const STableListInfo* pTableList, uint64_t uid, int32_t startIndex) {
9,886✔
3442
  int32_t numOfTables = taosArrayGetSize(pTableList->pTableList);
9,886✔
3443
  if (startIndex >= numOfTables) {
9,886✔
UNCOV
3444
    return -1;
×
3445
  }
3446

3447
  for (int32_t i = startIndex; i < numOfTables; ++i) {
91,619✔
3448
    STableKeyInfo* p = taosArrayGet(pTableList->pTableList, i);
91,619✔
3449
    if (!p) {
91,619✔
UNCOV
3450
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
UNCOV
3451
      return -1;
×
3452
    }
3453
    if (p->uid == uid) {
91,619✔
3454
      return i;
9,886✔
3455
    }
3456
  }
UNCOV
3457
  return -1;
×
3458
}
3459

3460
void tableListGetSourceTableInfo(const STableListInfo* pTableList, uint64_t* psuid, uint64_t* uid, int32_t* type) {
69,746✔
3461
  *psuid = pTableList->idInfo.suid;
69,746✔
3462
  *uid = pTableList->idInfo.uid;
69,746✔
3463
  *type = pTableList->idInfo.tableType;
69,746✔
3464
}
69,746✔
3465

3466
uint64_t tableListGetTableGroupId(const STableListInfo* pTableList, uint64_t tableUid) {
582,732,275✔
3467
  int32_t* slot = taosHashGet(pTableList->map, &tableUid, sizeof(tableUid));
582,732,275✔
3468
  if (slot == NULL) {
582,785,092✔
UNCOV
3469
    qDebug("table:%" PRIu64 " not found in table list", tableUid);
×
UNCOV
3470
    return -1;
×
3471
  }
3472

3473
  STableKeyInfo* pKeyInfo = taosArrayGet(pTableList->pTableList, *slot);
582,785,092✔
3474
  if (pKeyInfo == NULL) {
582,937,870✔
UNCOV
3475
    qDebug("table:%" PRIu64 " not found in table list", tableUid);
×
UNCOV
3476
    return -1;
×
3477
  }
3478
  return pKeyInfo->groupId;
582,937,870✔
3479
}
3480

3481
// TODO handle the group offset info, fix it, the rule of group output will be broken by this function
3482
// int32_t tableListRemoveTableInfo(STableListInfo* pTableList, uint64_t uid) {
3483
//   int32_t code = TSDB_CODE_SUCCESS;
3484
//   int32_t lino = 0;
3485

3486
//   int32_t* slot = taosHashGet(pTableList->map, &uid, sizeof(uid));
3487
//   if (slot == NULL) {
3488
//     qDebug("table:%" PRIu64 " not found in table list", uid);
3489
//     return 0;
3490
//   }
3491

3492
//   taosArrayRemove(pTableList->pTableList, *slot);
3493
//   code = taosHashRemove(pTableList->map, &uid, sizeof(uid));
3494

3495
//   _end:
3496
//   if (code != TSDB_CODE_SUCCESS) {
3497
//     qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
3498
//   } else {
3499
//     qDebug("uid:%" PRIu64 ", remove from table list", uid);
3500
//   }
3501

3502
//   return code;
3503
// }
3504

3505
int32_t tableListAddTableInfo(STableListInfo* pTableList, uint64_t uid, uint64_t gid) {
204,055✔
3506
  int32_t code = TSDB_CODE_SUCCESS;
204,055✔
3507
  int32_t lino = 0;
204,055✔
3508
  if (pTableList->map == NULL) {
204,055✔
UNCOV
3509
    pTableList->map = taosHashInit(32, taosGetDefaultHashFunction(TSDB_DATA_TYPE_BINARY), false, HASH_ENTRY_LOCK);
×
UNCOV
3510
    QUERY_CHECK_NULL(pTableList->map, code, lino, _end, terrno);
×
3511
  }
3512

3513
  STableKeyInfo keyInfo = {.uid = uid, .groupId = gid};
204,377✔
3514
  void*         p = taosHashGet(pTableList->map, &uid, sizeof(uid));
204,206✔
3515
  if (p != NULL) {
204,131✔
3516
    qInfo("table:%" PRId64 " already in tableIdList, ignore it", uid);
150✔
3517
    goto _end;
150✔
3518
  }
3519

3520
  void* tmp = taosArrayPush(pTableList->pTableList, &keyInfo);
203,981✔
3521
  QUERY_CHECK_NULL(tmp, code, lino, _end, terrno);
204,161✔
3522

3523
  int32_t slot = (int32_t)taosArrayGetSize(pTableList->pTableList) - 1;
204,161✔
3524
  code = taosHashPut(pTableList->map, &uid, sizeof(uid), &slot, sizeof(slot));
204,227✔
3525
  if (code != TSDB_CODE_SUCCESS) {
204,455✔
3526
    // we have checked the existence of uid in hash map above
UNCOV
3527
    QUERY_CHECK_CONDITION((code != TSDB_CODE_DUP_KEY), code, lino, _end, TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR);
×
UNCOV
3528
    taosArrayPopTailBatch(pTableList->pTableList, 1);  // let's pop the last element in the array list
×
3529
  }
3530

3531
_end:
204,605✔
3532
  if (code != TSDB_CODE_SUCCESS) {
204,434✔
UNCOV
3533
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
3534
  } else {
3535
    qDebug("uid:%" PRIu64 ", groupId:%" PRIu64 " added into table list, slot:%d, total:%d", uid, gid, slot, slot + 1);
204,434✔
3536
  }
3537

3538
  return code;
204,842✔
3539
}
3540

3541
int32_t tableListGetGroupList(const STableListInfo* pTableList, int32_t ordinalGroupIndex, STableKeyInfo** pKeyInfo,
189,240,882✔
3542
                              int32_t* size) {
3543
  int32_t totalGroups = tableListGetOutputGroups(pTableList);
189,240,882✔
3544
  int32_t numOfTables = 0;
189,252,445✔
3545
  int32_t code = tableListGetSize(pTableList, &numOfTables);
189,315,894✔
3546
  if (code != TSDB_CODE_SUCCESS) {
189,237,196✔
UNCOV
3547
    qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(code));
×
UNCOV
3548
    return code;
×
3549
  }
3550

3551
  if (ordinalGroupIndex < 0 || ordinalGroupIndex >= totalGroups) {
189,237,196✔
UNCOV
3552
    return TSDB_CODE_INVALID_PARA;
×
3553
  }
3554

3555
  // here handle two special cases:
3556
  // 1. only one group exists, and 2. one table exists for each group.
3557
  if (totalGroups == 1) {
189,237,196✔
3558
    *size = numOfTables;
188,825,169✔
3559
    *pKeyInfo = (*size == 0) ? NULL : taosArrayGet(pTableList->pTableList, 0);
188,912,145✔
3560
    return TSDB_CODE_SUCCESS;
188,722,679✔
3561
  } else if (totalGroups == numOfTables) {
412,057✔
3562
    *size = 1;
331,351✔
3563
    *pKeyInfo = taosArrayGet(pTableList->pTableList, ordinalGroupIndex);
331,351✔
3564
    return TSDB_CODE_SUCCESS;
331,351✔
3565
  }
3566

3567
  int32_t offset = pTableList->groupOffset[ordinalGroupIndex];
80,706✔
3568
  if (ordinalGroupIndex < totalGroups - 1) {
53,490✔
3569
    *size = pTableList->groupOffset[ordinalGroupIndex + 1] - offset;
40,428✔
3570
  } else {
3571
    *size = numOfTables - offset;
13,062✔
3572
  }
3573

3574
  *pKeyInfo = taosArrayGet(pTableList->pTableList, offset);
53,490✔
3575
  return TSDB_CODE_SUCCESS;
53,490✔
3576
}
3577

3578
int32_t tableListGetOutputGroups(const STableListInfo* pTableList) { return pTableList->numOfOuputGroups; }
564,540,396✔
3579

3580
bool oneTableForEachGroup(const STableListInfo* pTableList) { return pTableList->oneTableForEachGroup; }
484,317✔
3581

3582
STableListInfo* tableListCreate() {
225,112,273✔
3583
  STableListInfo* pListInfo = taosMemoryCalloc(1, sizeof(STableListInfo));
225,112,273✔
3584
  if (pListInfo == NULL) {
224,953,019✔
UNCOV
3585
    return NULL;
×
3586
  }
3587

3588
  pListInfo->remainGroups = NULL;
224,953,019✔
3589
  pListInfo->pTableList = taosArrayInit(4, sizeof(STableKeyInfo));
224,976,556✔
3590
  if (pListInfo->pTableList == NULL) {
225,015,044✔
UNCOV
3591
    goto _error;
×
3592
  }
3593

3594
  pListInfo->map = taosHashInit(1024, taosGetDefaultHashFunction(TSDB_DATA_TYPE_BIGINT), false, HASH_ENTRY_LOCK);
225,039,400✔
3595
  if (pListInfo->map == NULL) {
225,216,742✔
UNCOV
3596
    goto _error;
×
3597
  }
3598

3599
  pListInfo->numOfOuputGroups = 1;
225,210,891✔
3600
  return pListInfo;
225,214,677✔
3601

3602
_error:
×
UNCOV
3603
  tableListDestroy(pListInfo);
×
UNCOV
3604
  return NULL;
×
3605
}
3606

3607
void tableListDestroy(STableListInfo* pTableListInfo) {
234,177,427✔
3608
  if (pTableListInfo == NULL) {
234,177,427✔
3609
    return;
9,123,344✔
3610
  }
3611

3612
  taosArrayDestroy(pTableListInfo->pTableList);
225,054,083✔
3613
  taosMemoryFreeClear(pTableListInfo->groupOffset);
224,885,735✔
3614

3615
  taosHashCleanup(pTableListInfo->map);
224,946,034✔
3616
  taosHashCleanup(pTableListInfo->remainGroups);
225,138,441✔
3617
  pTableListInfo->pTableList = NULL;
225,151,362✔
3618
  pTableListInfo->map = NULL;
225,162,240✔
3619
  taosMemoryFree(pTableListInfo);
225,122,294✔
3620
}
3621

3622
void tableListClear(STableListInfo* pTableListInfo) {
151,708✔
3623
  if (pTableListInfo == NULL) {
151,708✔
UNCOV
3624
    return;
×
3625
  }
3626

3627
  taosArrayClear(pTableListInfo->pTableList);
151,708✔
3628
  taosHashClear(pTableListInfo->map);
151,804✔
3629
  taosHashClear(pTableListInfo->remainGroups);
152,002✔
3630
  taosMemoryFree(pTableListInfo->groupOffset);
152,002✔
3631
  pTableListInfo->numOfOuputGroups = 1;
152,002✔
3632
  pTableListInfo->oneTableForEachGroup = false;
152,002✔
3633
}
3634

3635
static int32_t orderbyGroupIdComparFn(const void* p1, const void* p2) {
488,238,257✔
3636
  STableKeyInfo* pInfo1 = (STableKeyInfo*)p1;
488,238,257✔
3637
  STableKeyInfo* pInfo2 = (STableKeyInfo*)p2;
488,238,257✔
3638

3639
  if (pInfo1->groupId == pInfo2->groupId) {
488,238,257✔
3640
    return 0;
472,076,097✔
3641
  } else {
3642
    return pInfo1->groupId < pInfo2->groupId ? -1 : 1;
16,162,955✔
3643
  }
3644
}
3645

3646
int32_t sortTableGroup(STableListInfo* pTableListInfo) {
18,765,809✔
3647
  int32_t code = TSDB_CODE_SUCCESS;
18,765,809✔
3648
  taosArraySort(pTableListInfo->pTableList, orderbyGroupIdComparFn);
18,765,809✔
3649
  int32_t size = taosArrayGetSize(pTableListInfo->pTableList);
18,769,895✔
3650
  if (size == 0) {
18,768,963✔
UNCOV
3651
    pTableListInfo->numOfOuputGroups = 0;
×
UNCOV
3652
    return code;
×
3653
  }
3654

3655
  SArray* pList = taosArrayInit(4, sizeof(int32_t));
18,768,963✔
3656
  if (!pList) {
18,768,809✔
UNCOV
3657
    code = terrno;
×
UNCOV
3658
    goto end;
×
3659
  }
3660

3661
  STableKeyInfo* pInfo = taosArrayGet(pTableListInfo->pTableList, 0);
18,768,809✔
3662
  if (pInfo == NULL) {
18,755,638✔
3663
    qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
UNCOV
3664
    code = terrno;
×
UNCOV
3665
    goto end;
×
3666
  }
3667
  uint64_t gid = pInfo->groupId;
18,755,638✔
3668

3669
  int32_t start = 0;
18,764,057✔
3670
  void*   tmp = taosArrayPush(pList, &start);
18,763,534✔
3671
  if (tmp == NULL) {
18,763,534✔
3672
    qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
UNCOV
3673
    code = terrno;
×
UNCOV
3674
    goto end;
×
3675
  }
3676

3677
  for (int32_t i = 1; i < size; ++i) {
120,054,059✔
3678
    pInfo = taosArrayGet(pTableListInfo->pTableList, i);
101,293,968✔
3679
    if (pInfo == NULL) {
101,293,151✔
3680
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
UNCOV
3681
      code = terrno;
×
UNCOV
3682
      goto end;
×
3683
    }
3684
    if (pInfo->groupId != gid) {
101,293,151✔
3685
      tmp = taosArrayPush(pList, &i);
3,516,665✔
3686
      if (tmp == NULL) {
3,516,665✔
3687
        qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
UNCOV
3688
        code = terrno;
×
UNCOV
3689
        goto end;
×
3690
      }
3691
      gid = pInfo->groupId;
3,516,665✔
3692
    }
3693
  }
3694

3695
  pTableListInfo->numOfOuputGroups = taosArrayGetSize(pList);
18,767,409✔
3696
  pTableListInfo->groupOffset = taosMemoryMalloc(sizeof(int32_t) * pTableListInfo->numOfOuputGroups);
18,769,061✔
3697
  if (pTableListInfo->groupOffset == NULL) {
18,768,776✔
UNCOV
3698
    code = terrno;
×
UNCOV
3699
    goto end;
×
3700
  }
3701

3702
  memcpy(pTableListInfo->groupOffset, taosArrayGet(pList, 0), sizeof(int32_t) * pTableListInfo->numOfOuputGroups);
18,757,681✔
3703

3704
end:
18,757,906✔
3705
  taosArrayDestroy(pList);
18,760,979✔
3706
  return code;
18,749,979✔
3707
}
3708

3709
int32_t buildGroupIdMapForAllTables(STableListInfo* pTableListInfo, SReadHandle* pHandle, SScanPhysiNode* pScanNode,
208,872,408✔
3710
                                    SNodeList* group, bool groupSort, uint8_t* digest, SStorageAPI* pAPI, SHashObj* groupIdMap) {
3711
  int32_t code = TSDB_CODE_SUCCESS;
208,872,408✔
3712

3713
  bool   groupByTbname = groupbyTbname(group);
208,872,408✔
3714
  size_t numOfTables = taosArrayGetSize(pTableListInfo->pTableList);
208,890,330✔
3715
  if (!numOfTables) {
208,828,467✔
3716
    return code;
6,214✔
3717
  }
3718
  qDebug("numOfTables:%zu, groupByTbname:%d, group:%p", numOfTables, groupByTbname, group);
208,822,253✔
3719
  if (group == NULL || groupByTbname) {
208,836,086✔
3720
    if (tsCountAlwaysReturnValue && QUERY_NODE_PHYSICAL_PLAN_TABLE_SCAN == nodeType(pScanNode) &&
206,965,650✔
3721
        ((STableScanPhysiNode*)pScanNode)->needCountEmptyTable) {
179,686,281✔
3722
      pTableListInfo->remainGroups =
7,213,297✔
3723
          taosHashInit(numOfTables, taosGetDefaultHashFunction(TSDB_DATA_TYPE_BIGINT), false, HASH_NO_LOCK);
7,213,297✔
3724
      if (pTableListInfo->remainGroups == NULL) {
7,213,297✔
UNCOV
3725
        return terrno;
×
3726
      }
3727

3728
      for (int i = 0; i < numOfTables; i++) {
25,605,682✔
3729
        STableKeyInfo* info = taosArrayGet(pTableListInfo->pTableList, i);
18,392,535✔
3730
        if (!info) {
18,392,385✔
UNCOV
3731
          qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
UNCOV
3732
          return terrno;
×
3733
        }
3734
        info->groupId = groupByTbname ? info->uid : 0;
18,392,385✔
3735
        int32_t tempRes = taosHashPut(pTableListInfo->remainGroups, &(info->groupId), sizeof(info->groupId),
18,392,535✔
3736
                                      &(info->uid), sizeof(info->uid));
18,392,535✔
3737
        if (tempRes != TSDB_CODE_SUCCESS && tempRes != TSDB_CODE_DUP_KEY) {
18,392,385✔
UNCOV
3738
          qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(tempRes));
×
UNCOV
3739
          return tempRes;
×
3740
        }
3741
      }
3742
    } else {
3743
      for (int32_t i = 0; i < numOfTables; i++) {
671,980,158✔
3744
        STableKeyInfo* info = taosArrayGet(pTableListInfo->pTableList, i);
472,275,865✔
3745
        if (!info) {
472,164,332✔
UNCOV
3746
          qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
UNCOV
3747
          return terrno;
×
3748
        }
3749
        info->groupId = groupByTbname ? info->uid : 0;
472,164,332✔
3750
        
3751
      }
3752
    }
3753
    if (groupIdMap && group != NULL){
206,917,440✔
3754
      getColInfoResultForGroupbyForStream(pHandle->vnode, group, pTableListInfo, pAPI, groupIdMap);
103,688✔
3755
    }
3756

3757
    pTableListInfo->oneTableForEachGroup = groupByTbname;
206,917,032✔
3758
    if (numOfTables == 1 && pTableListInfo->idInfo.tableType == TSDB_CHILD_TABLE) {
206,985,242✔
3759
      pTableListInfo->oneTableForEachGroup = true;
92,764,627✔
3760
    }
3761

3762
    if (groupSort && groupByTbname) {
206,969,032✔
3763
      taosArraySort(pTableListInfo->pTableList, orderbyGroupIdComparFn);
1,431,620✔
3764
      pTableListInfo->numOfOuputGroups = numOfTables;
1,431,212✔
3765
    } else if (groupByTbname && pScanNode->groupOrderScan) {
205,537,412✔
3766
      pTableListInfo->numOfOuputGroups = numOfTables;
30,754✔
3767
    } else {
3768
      pTableListInfo->numOfOuputGroups = 1;
205,506,658✔
3769
    }
3770
    if (groupSort || pScanNode->groupOrderScan) {
206,916,269✔
3771
      code = sortTableGroup(pTableListInfo);
18,588,289✔
3772
    }
3773
  } else {
3774
    bool initRemainGroups = false;
1,870,436✔
3775
    if (QUERY_NODE_PHYSICAL_PLAN_TABLE_SCAN == nodeType(pScanNode)) {
1,870,436✔
3776
      STableScanPhysiNode* pTableScanNode = (STableScanPhysiNode*)pScanNode;
1,706,339✔
3777
      if (tsCountAlwaysReturnValue && pTableScanNode->needCountEmptyTable &&
1,706,339✔
3778
          !(groupSort || pScanNode->groupOrderScan)) {
864,373✔
3779
        initRemainGroups = true;
837,191✔
3780
      }
3781
    }
3782

3783
    code = getColInfoResultForGroupby(pHandle->vnode, group, pTableListInfo, digest, pAPI, initRemainGroups, groupIdMap);
1,870,436✔
3784
    if (code != TSDB_CODE_SUCCESS) {
1,868,901✔
UNCOV
3785
      return code;
×
3786
    }
3787

3788
    if (pScanNode->groupOrderScan) pTableListInfo->numOfOuputGroups = taosArrayGetSize(pTableListInfo->pTableList);
1,868,901✔
3789

3790
    if (groupSort || pScanNode->groupOrderScan) {
1,869,746✔
3791
      code = sortTableGroup(pTableListInfo);
146,840✔
3792
    }
3793
  }
3794

3795
  // add all table entry in the hash map
3796
  size_t size = taosArrayGetSize(pTableListInfo->pTableList);
208,817,034✔
3797
  for (int32_t i = 0; i < size; ++i) {
708,743,682✔
3798
    STableKeyInfo* p = taosArrayGet(pTableListInfo->pTableList, i);
499,876,175✔
3799
    if (!p) {
499,702,342✔
UNCOV
3800
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(terrno));
×
UNCOV
3801
      return terrno;
×
3802
    }
3803
    int32_t tempRes = taosHashPut(pTableListInfo->map, &p->uid, sizeof(uint64_t), &i, sizeof(int32_t));
499,702,342✔
3804
    if (tempRes != TSDB_CODE_SUCCESS && tempRes != TSDB_CODE_DUP_KEY) {
499,899,982✔
UNCOV
3805
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(tempRes));
×
UNCOV
3806
      return tempRes;
×
3807
    }
3808
  }
3809

3810
  return code;
208,926,425✔
3811
}
3812

3813
int32_t createScanTableListInfo(SScanPhysiNode* pScanNode, SNodeList* pGroupTags, bool groupSort, SReadHandle* pHandle,
219,983,196✔
3814
                                STableListInfo* pTableListInfo, SNode* pTagCond, SNode* pTagIndexCond,
3815
                                SExecTaskInfo* pTaskInfo, SHashObj* groupIdMap) {
3816
  int64_t     st = taosGetTimestampUs();
220,016,254✔
3817
  const char* idStr = GET_TASKID(pTaskInfo);
220,016,254✔
3818

3819
  if (pHandle == NULL) {
219,732,375✔
UNCOV
3820
    qError("invalid handle, in creating operator tree, %s", idStr);
×
UNCOV
3821
    return TSDB_CODE_INVALID_PARA;
×
3822
  }
3823

3824
  if (pHandle->uid != 0) {
219,732,375✔
3825
    pScanNode->uid = pHandle->uid;
34,846✔
3826
    pScanNode->tableType = TSDB_CHILD_TABLE;
34,846✔
3827
  }
3828
  uint8_t digest[17] = {0};
219,974,810✔
3829
  int32_t code = getTableList(pHandle->vnode, pScanNode, pTagCond, pTagIndexCond, pTableListInfo, digest, idStr,
219,881,084✔
3830
                              &pTaskInfo->storageAPI, pTaskInfo->pStreamRuntimeInfo);
219,966,205✔
3831
  if (code != TSDB_CODE_SUCCESS) {
220,031,252✔
3832
    qError("failed to getTableList, code:%s", tstrerror(code));
1,108✔
3833
    return code;
1,108✔
3834
  }
3835

3836
  int32_t numOfTables = taosArrayGetSize(pTableListInfo->pTableList);
220,030,144✔
3837

3838
  int64_t st1 = taosGetTimestampUs();
220,080,376✔
3839
  pTaskInfo->cost.extractListTime = (st1 - st) / 1000.0;
220,080,376✔
3840
  qDebug("extract queried table list completed, %d tables, elapsed time:%.2f ms %s", numOfTables,
220,100,876✔
3841
         pTaskInfo->cost.extractListTime, idStr);
3842

3843
  if (numOfTables == 0) {
220,047,873✔
3844
    qDebug("no table qualified for query, %s", idStr);
11,269,358✔
3845
    return TSDB_CODE_SUCCESS;
11,269,358✔
3846
  }
3847

3848
  code = buildGroupIdMapForAllTables(pTableListInfo, pHandle, pScanNode, pGroupTags, groupSort, digest, &pTaskInfo->storageAPI, groupIdMap);
208,778,515✔
3849
  if (code != TSDB_CODE_SUCCESS) {
208,822,654✔
UNCOV
3850
    return code;
×
3851
  }
3852

3853
  pTaskInfo->cost.groupIdMapTime = (taosGetTimestampUs() - st1) / 1000.0;
208,826,245✔
3854
  qDebug("generate group id map completed, elapsed time:%.2f ms %s", pTaskInfo->cost.groupIdMapTime, idStr);
208,814,902✔
3855

3856
  return TSDB_CODE_SUCCESS;
208,776,098✔
3857
}
3858

3859
char* getStreamOpName(uint16_t opType) {
10,103,347✔
3860
  switch (opType) {
10,103,347✔
UNCOV
3861
    case QUERY_NODE_PHYSICAL_PLAN_STREAM_SCAN:
×
UNCOV
3862
      return "stream scan";
×
3863
    case QUERY_NODE_PHYSICAL_PLAN_PROJECT:
9,910,278✔
3864
      return "project";
9,910,278✔
3865
    case QUERY_NODE_PHYSICAL_PLAN_EXTERNAL_WINDOW:
193,945✔
3866
      return "external window";
193,945✔
3867
  }
UNCOV
3868
  return "error name";
×
3869
}
3870

3871
void printDataBlock(SSDataBlock* pBlock, const char* flag, const char* taskIdStr, int64_t qId) {
274,079,323✔
3872
  if (qDebugFlag & DEBUG_TRACE) {
274,079,323✔
3873
    if (!pBlock) {
38,352✔
3874
      qDebug("%" PRIx64 " %s %s %s: Block is Null", qId, taskIdStr, flag, __func__);
6,936✔
3875
      return;
6,936✔
3876
    } else if (pBlock->info.rows == 0) {
31,416✔
UNCOV
3877
      qDebug("%" PRIx64 " %s %s %s: Block is Empty. block type %d", qId, taskIdStr, flag, __func__, pBlock->info.type);
×
UNCOV
3878
      return;
×
3879
    }
3880
    
3881
    char*   pBuf = NULL;
31,416✔
3882
    int32_t code = dumpBlockData(pBlock, flag, &pBuf, taskIdStr, qId);
31,416✔
3883
    if (code == 0) {
31,416✔
3884
      qDebugL("%" PRIx64 " %s %s", qId, __func__, pBuf);
31,416✔
3885
      taosMemoryFree(pBuf);
31,416✔
3886
    }
3887
  }
3888
}
3889

3890
void printSpecDataBlock(SSDataBlock* pBlock, const char* flag, const char* opStr, const char* taskIdStr) {
×
3891
  if (!pBlock) {
×
3892
    qDebug("%s===stream===%s %s: Block is Null", taskIdStr, flag, opStr);
×
3893
    return;
×
UNCOV
3894
  } else if (pBlock->info.rows == 0) {
×
UNCOV
3895
    qDebug("%s===stream===%s %s: Block is Empty. block type %d.skey:%" PRId64 ",ekey:%" PRId64 ",version%" PRId64,
×
3896
           taskIdStr, flag, opStr, pBlock->info.type, pBlock->info.window.skey, pBlock->info.window.ekey,
3897
           pBlock->info.version);
3898
    return;
×
3899
  }
3900
  if (qDebugFlag & DEBUG_TRACE) {
×
3901
    char* pBuf = NULL;
×
3902
    char  flagBuf[64];
×
3903
    snprintf(flagBuf, sizeof(flagBuf), "%s %s", flag, opStr);
×
3904
    int32_t code = dumpBlockData(pBlock, flagBuf, &pBuf, taskIdStr, 0);
×
3905
    if (code == 0) {
×
UNCOV
3906
      qDebug("%s", pBuf);
×
UNCOV
3907
      taosMemoryFree(pBuf);
×
3908
    }
3909
  }
3910
}
3911

3912
TSKEY getStartTsKey(STimeWindow* win, const TSKEY* tsCols) { return tsCols == NULL ? win->skey : tsCols[0]; }
12,483,799✔
3913

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

3917
  int64_t duration = pWin->ekey > pWin->skey ? pWin->ekey - pWin->skey + delta : pWin->skey - pWin->ekey + delta;
2,147,483,647✔
3918
  ts[2] = duration;            // set the duration
2,147,483,647✔
3919
  ts[3] = pWin->skey;          // window start key
2,147,483,647✔
3920
  ts[4] = pWin->ekey + delta;  // window end key
2,147,483,647✔
3921
}
2,147,483,647✔
3922

3923
int32_t compKeys(const SArray* pSortGroupCols, const char* oldkeyBuf, int32_t oldKeysLen, const SSDataBlock* pBlock,
973,805,006✔
3924
                 int32_t rowIndex) {
3925
  SColumnDataAgg* pColAgg = NULL;
973,805,006✔
3926
  const char*     isNull = oldkeyBuf;
973,805,006✔
3927
  const char*     p = oldkeyBuf + sizeof(int8_t) * pSortGroupCols->size;
973,805,006✔
3928

3929
  for (int32_t i = 0; i < pSortGroupCols->size; ++i) {
2,147,483,647✔
3930
    const SColumn*         pCol = (SColumn*)TARRAY_GET_ELEM(pSortGroupCols, i);
1,509,795,782✔
3931
    const SColumnInfoData* pColInfoData = TARRAY_GET_ELEM(pBlock->pDataBlock, pCol->slotId);
1,511,360,759✔
3932
    if (pBlock->pBlockAgg) pColAgg = &pBlock->pBlockAgg[pCol->slotId];
1,511,410,654✔
3933

3934
    if (colDataIsNull(pColInfoData, pBlock->info.rows, rowIndex, pColAgg)) {
2,147,483,647✔
3935
      if (isNull[i] != 1) return 1;
103,336,269✔
3936
    } else {
3937
      if (isNull[i] != 0) return 1;
1,408,224,281✔
3938
      const char* val = colDataGetData(pColInfoData, rowIndex);
1,407,436,177✔
3939
      if (pCol->type == TSDB_DATA_TYPE_JSON) {
1,407,490,850✔
3940
        int32_t len = getJsonValueLen(val);
×
UNCOV
3941
        if (memcmp(p, val, len) != 0) return 1;
×
UNCOV
3942
        p += len;
×
3943
      } else if (IS_VAR_DATA_TYPE(pCol->type)) {
1,407,470,525✔
3944
        if (IS_STR_DATA_BLOB(pCol->type)) {
478,656,830✔
3945
          if (memcmp(p, val, blobDataTLen(val)) != 0) return 1;
165✔
UNCOV
3946
          p += blobDataTLen(val);
×
3947
        } else {
3948
          if (memcmp(p, val, varDataTLen(val)) != 0) return 1;
477,893,843✔
3949
          p += varDataTLen(val);
470,723,929✔
3950
        }
3951
      } else {
3952
        if (0 != memcmp(p, val, pCol->bytes)) return 1;
929,266,758✔
3953
        p += pCol->bytes;
912,009,359✔
3954
      }
3955
    }
3956
  }
3957
  if ((int32_t)(p - oldkeyBuf) != oldKeysLen) return 1;
948,681,039✔
3958
  return 0;
948,667,735✔
3959
}
3960

3961
int32_t buildKeys(char* keyBuf, const SArray* pSortGroupCols, const SSDataBlock* pBlock, int32_t rowIndex) {
25,534,670✔
3962
  uint32_t        colNum = pSortGroupCols->size;
25,534,670✔
3963
  SColumnDataAgg* pColAgg = NULL;
25,535,303✔
3964
  char*           isNull = keyBuf;
25,535,303✔
3965
  char*           p = keyBuf + sizeof(int8_t) * colNum;
25,535,303✔
3966

3967
  for (int32_t i = 0; i < colNum; ++i) {
76,289,662✔
3968
    const SColumn*         pCol = (SColumn*)TARRAY_GET_ELEM(pSortGroupCols, i);
50,752,249✔
3969
    const SColumnInfoData* pColInfoData = TARRAY_GET_ELEM(pBlock->pDataBlock, pCol->slotId);
50,753,304✔
3970
    if (pCol->slotId > pBlock->pDataBlock->size) continue;
50,753,726✔
3971

3972
    if (pBlock->pBlockAgg) pColAgg = &pBlock->pBlockAgg[pCol->slotId];
50,754,148✔
3973

3974
    if (colDataIsNull(pColInfoData, pBlock->info.rows, rowIndex, pColAgg)) {
101,497,113✔
3975
      isNull[i] = 1;
1,827,360✔
3976
    } else {
3977
      isNull[i] = 0;
48,924,045✔
3978
      const char* val = colDataGetData(pColInfoData, rowIndex);
48,921,302✔
3979
      if (pCol->type == TSDB_DATA_TYPE_JSON) {
48,924,889✔
3980
        int32_t len = getJsonValueLen(val);
×
UNCOV
3981
        memcpy(p, val, len);
×
UNCOV
3982
        p += len;
×
3983
      } else if (IS_VAR_DATA_TYPE(pCol->type)) {
48,915,183✔
3984
        if (IS_STR_DATA_BLOB(pCol->type)) {
7,347,071✔
UNCOV
3985
          blobDataCopy(p, val);
×
UNCOV
3986
          p += blobDataTLen(val);
×
3987
        } else {
3988
          varDataCopy(p, val);
7,341,585✔
3989
          p += varDataTLen(val);
7,341,585✔
3990
        }
3991
      } else {
3992
        memcpy(p, val, pCol->bytes);
41,580,983✔
3993
        p += pCol->bytes;
41,579,928✔
3994
      }
3995
    }
3996
  }
3997
  return (int32_t)(p - keyBuf);
25,537,413✔
3998
}
3999

4000
uint64_t calcGroupId(char* pData, int32_t len) {
2,147,483,647✔
4001
  T_MD5_CTX context;
2,147,483,647✔
4002
  tMD5Init(&context);
2,147,483,647✔
4003
  tMD5Update(&context, (uint8_t*)pData, len);
2,147,483,647✔
4004
  tMD5Final(&context);
2,147,483,647✔
4005

4006
  // NOTE: only extract the initial 8 bytes of the final MD5 digest
4007
  uint64_t id = 0;
2,147,483,647✔
4008
  memcpy(&id, context.digest, sizeof(uint64_t));
2,147,483,647✔
4009
  if (0 == id) memcpy(&id, context.digest + 8, sizeof(uint64_t));
2,147,483,647✔
4010
  return id;
2,147,483,647✔
4011
}
4012

4013
SNodeList* makeColsNodeArrFromSortKeys(SNodeList* pSortKeys) {
42,932✔
4014
  SNode*     node;
4015
  SNodeList* ret = NULL;
42,932✔
4016
  FOREACH(node, pSortKeys) {
130,800✔
4017
    SOrderByExprNode* pSortKey = (SOrderByExprNode*)node;
87,868✔
4018
    int32_t           code = nodesListMakeAppend(&ret, pSortKey->pExpr);
87,868✔
4019
    if (code != TSDB_CODE_SUCCESS) {
87,868✔
4020
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(code));
×
UNCOV
4021
      terrno = code;
×
UNCOV
4022
      return NULL;
×
4023
    }
4024
  }
4025
  return ret;
42,932✔
4026
}
4027

4028
int32_t extractKeysLen(const SArray* keys, int32_t* pLen) {
42,932✔
4029
  int32_t code = TSDB_CODE_SUCCESS;
42,932✔
4030
  int32_t lino = 0;
42,932✔
4031
  int32_t len = 0;
42,932✔
4032
  int32_t keyNum = taosArrayGetSize(keys);
42,932✔
4033
  for (int32_t i = 0; i < keyNum; ++i) {
109,440✔
4034
    SColumn* pCol = (SColumn*)taosArrayGet(keys, i);
66,508✔
4035
    QUERY_CHECK_NULL(pCol, code, lino, _end, terrno);
66,508✔
4036
    len += pCol->bytes;
66,508✔
4037
  }
4038
  len += sizeof(int8_t) * keyNum;  // null flag
42,932✔
4039
  *pLen = len;
42,932✔
4040

4041
_end:
42,932✔
4042
  if (code != TSDB_CODE_SUCCESS) {
42,932✔
UNCOV
4043
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
4044
  }
4045
  return code;
42,932✔
4046
}
4047

4048
int32_t parseErrorMsgFromAnalyticServer(SJson* pJson, const char* pId) {
×
4049
  int32_t code = TSDB_CODE_ANA_ANODE_RETURN_ERROR;
×
UNCOV
4050
  if (pJson == NULL) {
×
UNCOV
4051
    return code;
×
4052
  }
4053

UNCOV
4054
  char    pMsg[1024] = {0};
×
4055
  int32_t ret = tjsonGetStringValue(pJson, "msg", pMsg);
×
4056

4057
  if (ret == 0) {
×
4058
    qError("%s failed to exec imputation, msg:%s", pId, pMsg);
×
4059
    if (strstr(pMsg, "white noise") != NULL) {
×
4060
      code = TSDB_CODE_ANA_WN_DATA;
×
4061
    } else if (strstr(pMsg, "white-noise") != NULL) {
×
4062
      code = TSDB_CODE_ANA_WN_DATA;
×
UNCOV
4063
    } else if (strstr(pMsg, "[Errno 111] Connection refused") != NULL) {
×
UNCOV
4064
      code = TSDB_CODE_ANA_ALGO_NOT_LOAD;
×
4065
    }
4066
  } else {
UNCOV
4067
    qError("%s failed to extract msg from server, unknown error", pId);
×
4068
  }
4069

UNCOV
4070
  return code;
×
4071
}
4072

4073

4074
int32_t createBlockFromRemoteValueNode(SSDataBlock** ppBlock, SRemoteValueNode* pRemote) {
26,334,951✔
4075
  SValueNode* pVal = (SValueNode*)pRemote;
26,334,951✔
4076
  int32_t code = 0;
26,334,951✔
4077
  SSDataBlock* pBlock = taosMemoryCalloc(1, sizeof(SSDataBlock));
26,334,951✔
4078
  if (pBlock == NULL) {
26,332,796✔
UNCOV
4079
    return terrno;
×
4080
  }
4081

4082
  pBlock->pDataBlock = taosArrayInit(1, sizeof(SColumnInfoData));
26,332,796✔
4083
  if (pBlock->pDataBlock == NULL) {
26,333,314✔
4084
    code = terrno;
×
UNCOV
4085
    taosMemoryFree(pBlock);
×
UNCOV
4086
    return code;
×
4087
  }
4088

4089
  SColumnInfoData idata =
26,334,970✔
4090
      createColumnInfoData(pVal->node.resType.type, pVal->node.resType.bytes, 0);
26,335,569✔
4091
  idata.info.scale = pVal->node.resType.scale;
26,336,123✔
4092
  idata.info.precision = pVal->node.resType.precision;
26,336,123✔
4093

4094
  code = blockDataAppendColInfo(pBlock, &idata);
26,335,571✔
4095
  if (code != TSDB_CODE_SUCCESS) {
26,336,123✔
4096
    qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(code));
×
4097
    blockDataDestroy(pBlock);
×
UNCOV
4098
    *ppBlock = NULL;
×
UNCOV
4099
    return code;
×
4100
  }
4101

4102
  *ppBlock = pBlock;
26,336,123✔
4103

4104
  return code;
26,336,123✔
4105
}
4106

4107

4108
int32_t extractSingleRspBlock(SRetrieveTableRsp* pRetrieveRsp, SSDataBlock* pb) {
26,335,571✔
4109
  int32_t            code = TSDB_CODE_SUCCESS;
26,335,571✔
4110
  int32_t            lino = 0;
26,335,571✔
4111
  void*              decompBuf = NULL;
26,335,571✔
4112

4113
  char* pNextStart = pRetrieveRsp->data;
26,335,571✔
4114
  char* pStart = pNextStart;
26,335,019✔
4115

4116
  int32_t index = 0;
26,335,571✔
4117

4118
  if (pRetrieveRsp->compressed) {  // decompress the data
26,335,571✔
UNCOV
4119
    decompBuf = taosMemoryMalloc(pRetrieveRsp->payloadLen);
×
UNCOV
4120
    QUERY_CHECK_NULL(decompBuf, code, lino, _end, terrno);
×
4121
  }
4122

4123
  int32_t compLen = *(int32_t*)pStart;
26,334,420✔
4124
  pStart += sizeof(int32_t);
26,336,123✔
4125

4126
  int32_t rawLen = *(int32_t*)pStart;
26,335,571✔
4127
  pStart += sizeof(int32_t);
26,335,571✔
4128
  QUERY_CHECK_CONDITION((compLen <= rawLen && compLen != 0), code, lino, _end, TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR);
26,334,407✔
4129

4130
  pNextStart = pStart + compLen;
26,334,407✔
4131
  if (pRetrieveRsp->compressed && (compLen < rawLen)) {
26,335,524✔
4132
    int32_t t = tsDecompressString(pStart, compLen, 1, decompBuf, rawLen, ONE_STAGE_COMP, NULL, 0);
×
UNCOV
4133
    QUERY_CHECK_CONDITION((t == rawLen), code, lino, _end, TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR);
×
UNCOV
4134
    pStart = decompBuf;
×
4135
  }
4136

4137
  code = blockDecodeInternal(pb, pStart, (const char**)&pStart);
26,334,407✔
4138
  if (code != 0) {
26,333,866✔
UNCOV
4139
    taosMemoryFreeClear(pRetrieveRsp);
×
UNCOV
4140
    goto _end;
×
4141
  }
4142

4143
_end:
26,333,866✔
4144
  if (code != TSDB_CODE_SUCCESS) {
26,334,970✔
UNCOV
4145
    blockDataDestroy(pb);
×
UNCOV
4146
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
4147
  }
4148
  return code;
26,334,970✔
4149
}
4150

4151
int32_t setValueFromResBlock(STaskSubJobCtx* ctx, SRemoteValueNode* pRes, SSDataBlock* pBlock) {
26,335,569✔
4152
  int32_t code = 0;
26,335,569✔
4153
  bool needFree = true;
26,335,569✔
4154
  int32_t colNum = taosArrayGetSize(pBlock->pDataBlock);
26,335,569✔
4155
  if (NULL == pBlock->pDataBlock || 1 != colNum || pBlock->info.rows > 1) {
26,332,184✔
4156
    qError("%s invalid scl fetch res block, pDataBlock:%p, colNum:%d, rows:%" PRId64, 
1,110✔
4157
      ctx->idStr, pBlock->pDataBlock, colNum, pBlock->info.rows);
UNCOV
4158
    return TSDB_CODE_PAR_INVALID_SCALAR_SUBQ_RES_ROWS;
×
4159
  }
4160
  
4161
  pRes->val.node.type = QUERY_NODE_VALUE;
26,330,524✔
4162
  pRes->val.flag &= (~VALUE_FLAG_VAL_UNSET);
26,330,526✔
4163
  pRes->val.translate = true;
26,333,303✔
4164
  
4165
  SColumnInfoData* pCol = taosArrayGet(pBlock->pDataBlock, 0);
26,331,690✔
4166
  if (colDataIsNull_s(pCol, 0)) {
26,333,902✔
4167
    pRes->val.isNull = true;
2,163,521✔
4168
  } else {
4169
    code = nodesSetValueNodeValueExt(&pRes->val, colDataGetData(pCol, 0), &needFree);
24,170,381✔
4170
  }
4171

4172
  if (!needFree) {
26,334,418✔
4173
    pCol->pData = NULL;
17,304✔
4174
  }
4175

4176
  return code;
26,334,418✔
4177
}
4178

4179
int32_t remoteFetchCallBack(void* param, SDataBuf* pMsg, int32_t code) {
47,620,723✔
4180
  SScalarFetchParam* pParam = (SScalarFetchParam*)param;
47,620,723✔
4181
  STaskSubJobCtx* ctx = pParam->pSubJobCtx;
47,620,723✔
4182
  SSDataBlock* pResBlock = NULL;
47,621,880✔
4183
  
4184
  taosMemoryFreeClear(pMsg->pEpSet);
47,620,682✔
4185

4186
  if (NULL == ctx) {
47,620,723✔
4187
    qWarn("scl fetch ctx not exists since it may have been released");
7,819✔
4188
    goto _exit;
7,819✔
4189
  }
4190

4191
  qDebug("%s subQIdx %d got rsp, code:%d, rsp:%p", ctx->idStr, pParam->subQIdx, code, pMsg->pData);
47,612,904✔
4192

4193
  taosWLockLatch(&ctx->lock);
47,612,904✔
4194
  ctx->param = NULL;
47,611,150✔
4195
  taosWUnLockLatch(&ctx->lock);
47,612,307✔
4196

4197
  if (ctx->transporterId > 0) {
47,612,859✔
4198
    int32_t ret = asyncFreeConnById(ctx->rpcHandle, ctx->transporterId);
47,611,751✔
4199
    if (ret != 0) {
47,614,061✔
UNCOV
4200
      qDebug("%s failed to free subQ rpc handle, code:%s, subQIdx:%d", ctx->idStr, tstrerror(ret), pParam->subQIdx);
×
4201
    }
4202
    ctx->transporterId = -1;
47,614,061✔
4203
  }
4204

4205
  if (0 == code && NULL == pMsg->pData) {
47,615,169✔
UNCOV
4206
    qError("%s invalid rsp msg, msgType:%d, len:%d", ctx->idStr, pMsg->msgType, pMsg->len);
×
UNCOV
4207
    code = TSDB_CODE_QRY_INVALID_MSG;
×
4208
  }
4209

4210
  if (code == TSDB_CODE_SUCCESS) {
47,612,904✔
4211
    SRetrieveTableRsp* pRsp = pMsg->pData;
39,331,838✔
4212
    pRsp->numOfRows = htobe64(pRsp->numOfRows);
39,332,392✔
4213
    pRsp->compLen = htonl(pRsp->compLen);
39,330,169✔
4214
    pRsp->payloadLen = htonl(pRsp->payloadLen);
39,331,286✔
4215
    pRsp->numOfCols = htonl(pRsp->numOfCols);
39,330,135✔
4216
    pRsp->useconds = htobe64(pRsp->useconds);
39,331,284✔
4217
    pRsp->numOfBlocks = htonl(pRsp->numOfBlocks);
39,330,182✔
4218

4219
    if (pRsp->numOfRows > 1 || pRsp->numOfBlocks > 1 || !pRsp->completed) {
39,326,252✔
4220
      qError("%s invalid scl fetch rsp received, subQIdx:%d, rows:%" PRId64 ", blocks:%d, completed:%d", 
912,233✔
4221
        ctx->idStr, pParam->subQIdx, pRsp->numOfRows, pRsp->numOfBlocks, pRsp->completed);
4222
      ctx->code = TSDB_CODE_PAR_INVALID_SCALAR_SUBQ_RES_ROWS;
912,233✔
4223
    } else if (0 == pRsp->numOfRows) {
38,418,497✔
4224
      SRemoteValueNode* pRemote = (SRemoteValueNode*)pParam->pRes;
12,086,813✔
4225
      pRemote->val.node.type = QUERY_NODE_VALUE;
12,086,813✔
4226
      pRemote->val.isNull = true;
12,087,365✔
4227
      pRemote->val.translate = true;
12,086,259✔
4228
      pRemote->val.flag &= (~VALUE_FLAG_VAL_UNSET);
12,086,259✔
4229
      taosArraySet(ctx->subResValues, pParam->subQIdx, &pParam->pRes);
12,086,811✔
4230
    } else {
4231
      qDebug("%s scl fetch rsp received, subQIdx:%d, rows:%" PRId64 , ctx->idStr, pParam->subQIdx, pRsp->numOfRows);
26,335,569✔
4232
      ctx->code = createBlockFromRemoteValueNode(&pResBlock, pParam->pRes);
26,335,569✔
4233
      if (TSDB_CODE_SUCCESS == ctx->code) {
26,334,463✔
4234
        ctx->code = blockDataEnsureCapacity(pResBlock, 1);
26,334,418✔
4235
      }
4236
      if (TSDB_CODE_SUCCESS == ctx->code) {
26,336,123✔
4237
        ctx->code = extractSingleRspBlock(pRsp, pResBlock);
26,334,972✔
4238
      }
4239
      if (TSDB_CODE_SUCCESS == ctx->code) {
26,336,155✔
4240
        ctx->code = setValueFromResBlock(ctx, pParam->pRes, pResBlock);
26,336,123✔
4241
      }
4242
      if (TSDB_CODE_SUCCESS == ctx->code) {
26,333,898✔
4243
        taosArraySet(ctx->subResValues, pParam->subQIdx, &pParam->pRes);
26,333,900✔
4244
      }
4245
    }
4246
  } else {
4247
    ctx->code = rpcCvtErrCode(code);
8,281,066✔
4248
    if (ctx->code != code) {
8,281,669✔
UNCOV
4249
      qError("%s scl fetch rsp received, subQIdx:%d, error:%s, cvted error: %s", ctx->idStr, pParam->subQIdx,
×
4250
             tstrerror(code), tstrerror(ctx->code));
4251
    } else {
4252
      qError("%s scl fetch rsp received, subQIdx:%d, error:%s", ctx->idStr, pParam->subQIdx, tstrerror(code));
8,281,066✔
4253
    }
4254
  }
4255
  
4256
  code = tsem_post(&pParam->pSubJobCtx->ready);
47,612,942✔
4257
  if (code != TSDB_CODE_SUCCESS) {
47,613,462✔
UNCOV
4258
    qError("failed to invoke post when scl fetch rsp is ready, code:%s", tstrerror(code));
×
4259
  }
4260

4261
_exit:
47,621,281✔
4262

4263
  taosMemoryFree(pMsg->pData);
47,620,727✔
4264
  blockDataDestroy(pResBlock);
47,620,177✔
4265

4266
  return code;
47,620,774✔
4267
}
4268

4269

4270
int32_t fetchRemoteValueImpl(STaskSubJobCtx* ctx, int32_t subQIdx, SRemoteValueNode* pRes) {
47,608,259✔
4271
  int32_t          code = TSDB_CODE_SUCCESS;
47,608,259✔
4272
  int32_t          lino = 0;
47,608,259✔
4273
  SDownstreamSourceNode* pSource = (SDownstreamSourceNode*)taosArrayGetP(ctx->subEndPoints, subQIdx);
47,608,259✔
4274

4275
  SResFetchReq req = {0};
47,611,828✔
4276
  req.header.vgId = pSource->addr.nodeId;
47,611,828✔
4277
  req.sId = pSource->sId;
47,612,291✔
4278
  req.clientId = pSource->clientId;
47,614,205✔
4279
  req.taskId = pSource->taskId;
47,597,551✔
4280
  req.queryId = ctx->queryId;
47,597,092✔
4281
  req.execId = pSource->execId;
47,590,728✔
4282

4283
  int32_t msgSize = tSerializeSResFetchReq(NULL, 0, &req, false);
47,600,948✔
4284
  if (msgSize < 0) {
47,605,545✔
UNCOV
4285
    return msgSize;
×
4286
  }
4287

4288
  void* msg = taosMemoryCalloc(1, msgSize);
47,605,545✔
4289
  if (NULL == msg) {
47,576,871✔
UNCOV
4290
    return terrno;
×
4291
  }
4292

4293
  msgSize = tSerializeSResFetchReq(msg, msgSize, &req, false);
47,576,871✔
4294
  if (msgSize < 0) {
47,603,262✔
UNCOV
4295
    taosMemoryFree(msg);
×
UNCOV
4296
    return msgSize;
×
4297
  }
4298

4299
  qDebug("%s scl build fetch msg and send to nodeId:%d, ep:%s, clientId:0x%" PRIx64 " taskId:0x%" PRIx64
47,603,262✔
4300
         ", execId:%d",
4301
         ctx->idStr, pSource->addr.nodeId, pSource->addr.epSet.eps[0].fqdn, pSource->clientId,
4302
         pSource->taskId, pSource->execId);
4303

4304
  // send the fetch remote task result reques
4305
  SMsgSendInfo* pMsgSendInfo = taosMemoryCalloc(1, sizeof(SMsgSendInfo));
47,606,639✔
4306
  if (NULL == pMsgSendInfo) {
47,571,709✔
4307
    taosMemoryFreeClear(msg);
×
UNCOV
4308
    qError("%s prepare message %d failed", ctx->idStr, (int32_t)sizeof(SMsgSendInfo));
×
UNCOV
4309
    return terrno;
×
4310
  }
4311

4312
  SScalarFetchParam* param = taosMemoryMalloc(sizeof(SScalarFetchParam));
47,571,709✔
4313
  if (NULL == param) {
47,597,139✔
4314
    taosMemoryFreeClear(msg);
×
4315
    taosMemoryFreeClear(pMsgSendInfo);
×
UNCOV
4316
    qError("%s prepare param %d failed", ctx->idStr, (int32_t)sizeof(SScalarFetchParam));
×
UNCOV
4317
    return terrno;
×
4318
  }
4319

4320
  taosWLockLatch(&ctx->lock);
47,597,139✔
4321
  
4322
  if (ctx->code) {
47,617,394✔
4323
    qError("task has been killed, error:%s", tstrerror(ctx->code));
×
4324
    taosMemoryFree(param);
×
UNCOV
4325
    code = ctx->code;
×
UNCOV
4326
    goto _end;
×
4327
  } else {
4328
    ctx->param = param;
47,603,907✔
4329
  }
4330
  
4331
  taosWUnLockLatch(&ctx->lock);
47,611,266✔
4332

4333
  param->subQIdx = subQIdx;
47,612,458✔
4334
  param->pRes = pRes;
47,613,048✔
4335
  param->pSubJobCtx = ctx;
47,599,370✔
4336

4337
  pMsgSendInfo->param = param;
47,610,756✔
4338
  pMsgSendInfo->paramFreeFp = taosAutoMemoryFree;
47,593,697✔
4339
  pMsgSendInfo->msgInfo.pData = msg;
47,597,042✔
4340
  pMsgSendInfo->msgInfo.len = msgSize;
47,601,016✔
4341
  pMsgSendInfo->msgType = pSource->fetchMsgType;
47,599,594✔
4342
  pMsgSendInfo->fp = remoteFetchCallBack;
47,571,680✔
4343
  pMsgSendInfo->requestId = ctx->queryId;
47,594,873✔
4344

4345
  code = asyncSendMsgToServer(ctx->rpcHandle, &pSource->addr.epSet, &ctx->transporterId, pMsgSendInfo);
47,616,916✔
4346
  QUERY_CHECK_CODE(code, lino, _end);
47,621,336✔
4347

4348
  code = qSemWait(ctx->pTaskInfo, &ctx->ready);
47,621,336✔
4349
  if (isTaskKilled(ctx->pTaskInfo)) {
47,622,029✔
4350
    code = getTaskCode(ctx->pTaskInfo);
13,849✔
4351
  } else {
4352
    code = ctx->code;
47,609,275✔
4353
  }
4354
      
4355
_end:
47,622,583✔
4356

4357
  taosWLockLatch(&ctx->lock);
47,622,583✔
4358
  ctx->param = NULL;
47,623,137✔
4359
  taosWUnLockLatch(&ctx->lock);
47,623,689✔
4360

4361
  if (code != TSDB_CODE_SUCCESS) {
47,623,689✔
4362
    qError("%s %s failed at line %d since %s", ctx->idStr, __func__, lino, tstrerror(code));
9,200,201✔
4363
  }
4364
  return code;
47,622,029✔
4365
}
4366

4367

4368
int32_t qFetchRemoteValue(void* pCtx, int32_t subQIdx, SRemoteValueNode* pRes) {
49,191,487✔
4369
  STaskSubJobCtx*  ctx = (STaskSubJobCtx*)pCtx;
49,191,487✔
4370
  int32_t code = 0, lino = 0;
49,191,487✔
4371
  int32_t       subEndPoinsNum = taosArrayGetSize(ctx->subEndPoints);
49,191,487✔
4372
  if (subQIdx >= subEndPoinsNum) {
49,193,873✔
UNCOV
4373
    qError("%s invalid subQIdx %d, subEndPointsNum:%d", ctx->idStr, subQIdx, subEndPoinsNum);
×
UNCOV
4374
    return TSDB_CODE_QRY_SUBQ_NOT_FOUND;
×
4375
  }
4376

4377
  SValueNode** ppRes = taosArrayGet(ctx->subResValues, subQIdx);
49,193,873✔
4378
  if (NULL == *ppRes) {
49,196,452✔
4379
    TAOS_CHECK_EXIT(fetchRemoteValueImpl(ctx, subQIdx, pRes));
47,613,824✔
4380
    *ppRes = (SValueNode*)pRes;
38,421,276✔
4381
  } else {
4382
    TAOS_CHECK_EXIT(valueNodeCopy(*ppRes, &pRes->val));
1,584,320✔
4383
    pRes->val.node.type = QUERY_NODE_VALUE;
1,584,320✔
4384
  }
4385

4386
_exit:
49,206,351✔
4387

4388
  if (code) {
49,206,351✔
4389
    qError("%s %s failed at line %d since %s", ctx->idStr, __func__, lino, tstrerror(code));
9,200,201✔
4390
  }
4391

4392
  return code;
49,205,797✔
4393
}
4394

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