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

taosdata / TDengine / #5034

24 Apr 2026 11:25AM UTC coverage: 73.058%. Remained the same
#5034

push

travis-ci

web-flow
merge: from main to 3.0 branch #35224

merge: from main to 3.0 branch[manual-only]

1336 of 1975 new or added lines in 48 files covered. (67.65%)

14149 existing lines in 164 files now uncovered.

275896 of 377640 relevant lines covered (73.06%)

132944440.29 hits per line

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

90.13
/source/libs/executor/src/timesliceoperator.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
#include "executorInt.h"
16
#include "filter.h"
17
#include "function.h"
18
#include "functionMgt.h"
19
#include "operator.h"
20
#include "querytask.h"
21
#include "storageapi.h"
22
#include "tcommon.h"
23
#include "tcompare.h"
24
#include "tdatablock.h"
25
#include "tfill.h"
26
#include "ttime.h"
27

28
typedef struct STimeSliceOperatorInfo {
29
  SSDataBlock*         pRes;
30
  STimeWindow          win;
31
  SNode*               pWin;        // for stream
32
  SInterval            interval;
33
  int64_t              current;
34
  SArray*              pPrevRow;     // SArray<SGroupValue>
35
  SArray*              pNextRow;     // SArray<SGroupValue>
36
  SArray*              pLinearInfo;  // SArray<SFillLinearInfo>
37
  bool                 isPrevRowSet;
38
  bool                 isNextRowSet;
39
  int32_t              fillType;      // fill type
40
  SColumn              tsCol;         // primary timestamp column
41
  SExprSupp            scalarSup;     // scalar calculation
42
  struct SFillColInfo* pFillColInfo;  // fill column info
43
  SRowKey              prevKey;       // record previous row key
44
  bool                 prevTsSet;     // denotes if previous timestamp is set
45
  uint64_t             groupId;
46
  SArray*              pPrevGroupKeys;
47
  SSDataBlock*         pNextGroupRes;
48
  SSDataBlock*         pRemainRes;   // save block unfinished processing
49
  int32_t              remainIndex;  // the remaining index in the block to be processed
50
  bool                 hasPk;
51
  SColumn              pkCol;
52
  bool                 prevNotified;
53
  bool                 nextNotified;
54
  int64_t              surroundingTime;
55
} STimeSliceOperatorInfo;
56

57
static void destroyTimeSliceOperatorInfo(void* param);
58

59
static void doKeepPrevRows(STimeSliceOperatorInfo* pSliceInfo, const SSDataBlock* pBlock, int32_t rowIndex) {
2,147,483,647✔
60
  int32_t numOfCols = taosArrayGetSize(pBlock->pDataBlock);
2,147,483,647✔
61
  for (int32_t i = 0; i < numOfCols; ++i) {
2,147,483,647✔
62
    SColumnInfoData* pColInfoData = taosArrayGet(pBlock->pDataBlock, i);
2,147,483,647✔
63

64
    SGroupKeys* pkey = taosArrayGet(pSliceInfo->pPrevRow, i);
2,147,483,647✔
65
    if (!colDataIsNull_s(pColInfoData, rowIndex)) {
2,147,483,647✔
66
      pkey->isNull = false;
2,147,483,647✔
67
      char* val = colDataGetData(pColInfoData, rowIndex);
2,147,483,647✔
68
      if (IS_VAR_DATA_TYPE(pkey->type)) {
2,147,483,647✔
69
        int32_t bytes = calcStrBytesByType(pkey->type, val);
1,058,915,307✔
70
        memcpy(pkey->pData, val, bytes);
1,062,038,829✔
71
      } else {
72
        memcpy(pkey->pData, val, pkey->bytes);
2,147,483,647✔
73
      }
74
    } else {
75
      pkey->isNull = true;
521,445,581✔
76
    }
77
  }
78

79
  pSliceInfo->isPrevRowSet = true;
2,147,483,647✔
80
}
2,147,483,647✔
81

82
static void doKeepNextRows(STimeSliceOperatorInfo* pSliceInfo, const SSDataBlock* pBlock, int32_t rowIndex) {
263,193,017✔
83
  int32_t numOfCols = taosArrayGetSize(pBlock->pDataBlock);
263,193,017✔
84
  for (int32_t i = 0; i < numOfCols; ++i) {
1,054,433,954✔
85
    SColumnInfoData* pColInfoData = taosArrayGet(pBlock->pDataBlock, i);
791,206,023✔
86

87
    SGroupKeys* pkey = taosArrayGet(pSliceInfo->pNextRow, i);
791,201,452✔
88
    if (!colDataIsNull_s(pColInfoData, rowIndex)) {
1,582,336,795✔
89
      pkey->isNull = false;
790,017,332✔
90
      char* val = colDataGetData(pColInfoData, rowIndex);
790,019,645✔
91
      if (!IS_VAR_DATA_TYPE(pkey->type)) {
790,084,882✔
92
        memcpy(pkey->pData, val, pkey->bytes);
789,475,353✔
93
      } else {
94
        int32_t bytes = calcStrBytesByType(pkey->type, val);
83,351✔
95
        memcpy(pkey->pData, val, bytes);
607,920✔
96
      }
97
    } else {
98
      pkey->isNull = true;
1,160,113✔
99
    }
100
  }
101

102
  pSliceInfo->isNextRowSet = true;
263,227,931✔
103
}
263,233,794✔
104

105
static void doKeepLinearInfo(STimeSliceOperatorInfo* pSliceInfo, const SSDataBlock* pBlock, int32_t rowIndex) {
2,147,483,647✔
106
  int32_t numOfCols = taosArrayGetSize(pBlock->pDataBlock);
2,147,483,647✔
107
  for (int32_t i = 0; i < numOfCols; ++i) {
2,147,483,647✔
108
    SColumnInfoData* pColInfoData = taosArrayGet(pBlock->pDataBlock, i);
2,147,483,647✔
109
    SColumnInfoData* pTsCol = taosArrayGet(pBlock->pDataBlock, pSliceInfo->tsCol.slotId);
2,147,483,647✔
110
    SFillLinearInfo* pLinearInfo = taosArrayGet(pSliceInfo->pLinearInfo, i);
2,147,483,647✔
111

112
    if (!IS_MATHABLE_TYPE(pColInfoData->info.type)) {
2,147,483,647✔
113
      continue;
1,062,038,829✔
114
    }
115

116
    // null value is represented by using key = INT64_MIN for now.
117
    // TODO: optimize to ignore null values for linear interpolation.
118
    if (!pLinearInfo->isStartSet) {
2,147,483,647✔
119
      if (!colDataIsNull_s(pColInfoData, rowIndex)) {
22,234,342✔
120
        pLinearInfo->start.key = *(int64_t*)colDataGetData(pTsCol, rowIndex);
11,082,564✔
121
        char* p = colDataGetData(pColInfoData, rowIndex);
11,082,564✔
122
        if (IS_VAR_DATA_TYPE(pColInfoData->info.type)) {
11,082,564✔
123
          if (IS_STR_DATA_BLOB(pColInfoData->info.type)) {
×
124
            memcpy(pLinearInfo->start.val, p, blobDataTLen(p));
×
125
          } else {
126
            memcpy(pLinearInfo->start.val, p, varDataTLen(p));
×
127
          }
128
        } else {
129
          memcpy(pLinearInfo->start.val, p, pLinearInfo->bytes);
11,082,564✔
130
        }
131
      }
132
      pLinearInfo->isStartSet = true;
11,117,171✔
133
    } else if (!pLinearInfo->isEndSet) {
2,147,483,647✔
134
      if (!colDataIsNull_s(pColInfoData, rowIndex)) {
15,256,328✔
135
        pLinearInfo->end.key = *(int64_t*)colDataGetData(pTsCol, rowIndex);
7,578,600✔
136

137
        char* p = colDataGetData(pColInfoData, rowIndex);
7,578,600✔
138
        if (IS_VAR_DATA_TYPE(pColInfoData->info.type)) {
7,578,600✔
139
          if (IS_STR_DATA_BLOB(pColInfoData->info.type)) {
×
140
            memcpy(pLinearInfo->end.val, p, blobDataTLen(p));
×
141
          } else {
142
            memcpy(pLinearInfo->end.val, p, varDataTLen(p));
×
143
          }
144
        } else {
145
          memcpy(pLinearInfo->end.val, p, pLinearInfo->bytes);
7,578,600✔
146
        }
147
      }
148
      pLinearInfo->isEndSet = true;
7,628,164✔
149
    } else {
150
      pLinearInfo->start.key = pLinearInfo->end.key;
2,147,483,647✔
151
      memcpy(pLinearInfo->start.val, pLinearInfo->end.val, pLinearInfo->bytes);
2,147,483,647✔
152

153
      if (!colDataIsNull_s(pColInfoData, rowIndex)) {
2,147,483,647✔
154
        pLinearInfo->end.key = *(int64_t*)colDataGetData(pTsCol, rowIndex);
2,147,483,647✔
155

156
        char* p = colDataGetData(pColInfoData, rowIndex);
2,147,483,647✔
157
        if (IS_VAR_DATA_TYPE(pColInfoData->info.type)) {
2,147,483,647✔
158
          if (IS_STR_DATA_BLOB(pColInfoData->info.type)) {
8,787,354✔
159
            memcpy(pLinearInfo->end.val, p, blobDataTLen(p));
×
160
          } else {
161
            memcpy(pLinearInfo->end.val, p, varDataTLen(p));
×
162
          }
163
        } else {
164
          memcpy(pLinearInfo->end.val, p, pLinearInfo->bytes);
2,147,483,647✔
165
        }
166

167
      } else {
168
        pLinearInfo->end.key = INT64_MIN;
511,425,526✔
169
      }
170
    }
171
  }
172
}
2,147,483,647✔
173

174
static FORCE_INLINE int32_t timeSliceEnsureBlockCapacity(STimeSliceOperatorInfo* pSliceInfo, SSDataBlock* pBlock) {
175
  if (pBlock->info.rows < pBlock->info.capacity) {
699,771,573✔
176
    return TSDB_CODE_SUCCESS;
699,767,488✔
177
  }
178

179
  uint32_t winNum = (pSliceInfo->win.ekey - pSliceInfo->win.skey) / pSliceInfo->interval.interval;
3,464✔
180
  uint32_t newRowsNum = pBlock->info.rows + TMIN(winNum / 4 + 1, 1048576);
3,464✔
181
  int32_t  code = blockDataEnsureCapacity(pBlock, newRowsNum);
3,464✔
182
  if (code != TSDB_CODE_SUCCESS) {
3,464✔
183
    qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(code));
×
184
    return code;
×
185
  }
186

187
  return TSDB_CODE_SUCCESS;
3,464✔
188
}
189

190
bool isIrowtsPseudoColumn(SExprInfo* pExprInfo) {
2,147,483,647✔
191
  char* name = pExprInfo->pExpr->_function.functionName;
2,147,483,647✔
192
  return (IS_TIMESTAMP_TYPE(pExprInfo->base.resSchema.type) && strcasecmp(name, "_irowts") == 0);
2,147,483,647✔
193
}
194

195
bool isIsfilledPseudoColumn(SExprInfo* pExprInfo) {
2,147,483,647✔
196
  char* name = pExprInfo->pExpr->_function.functionName;
2,147,483,647✔
197
  return (IS_BOOLEAN_TYPE(pExprInfo->base.resSchema.type) && strcasecmp(name, "_isfilled") == 0);
2,147,483,647✔
198
}
199

200
bool isIrowtsOriginPseudoColumn(SExprInfo* pExprInfo) {
1,924,301,822✔
201
  const char* name = pExprInfo->pExpr->_function.functionName;
1,924,301,822✔
202
  return (IS_TIMESTAMP_TYPE(pExprInfo->base.resSchema.type) && strcasecmp(name, "_irowts_origin") == 0);
1,924,307,509✔
203
}
204

205
static void tRowGetKeyFromColData(int64_t ts, SColumnInfoData* pPkCol, int32_t rowIndex, SRowKey* pKey) {
×
206
  pKey->ts = ts;
×
207
  pKey->numOfPKs = 1;
×
208

209
  int8_t t = pPkCol->info.type;
×
210

211
  pKey->pks[0].type = t;
×
212
  if (IS_NUMERIC_TYPE(t)) {
×
213
    valueSetDatum(pKey->pks, t, colDataGetData(pPkCol, rowIndex), tDataTypes[t].bytes);
×
214
  } else {
215
    char* p = colDataGetVarData(pPkCol, rowIndex);
×
216
    pKey->pks[0].pData = (uint8_t*)varDataVal(p);
×
217
    pKey->pks[0].nData = varDataLen(p);
×
218
  }
219
}
×
220

221
typedef enum {
222
  INVALID_TIMESTAMP_REASON_NONE = 0,  /* not invalid */
223
  INVALID_TIMESTAMP_REASON_PREV_TS_EQUAL = 1,
224
  INVALID_TIMESTAMP_REASON_PREV_TS_SMALLER = 2,
225
} EInvalidTimestampReason;
226

227
/**
228
  @brief Timestamp is invalid if current timestamp <= previous timestamp.
229
  Only timestamp is considered even if composite primary key exists.
230
*/
231
static EInvalidTimestampReason isInvalidTimestamp(
2,147,483,647✔
232
  STimeSliceOperatorInfo* pSliceInfo, int64_t currentTs,
233
  SColumnInfoData* pPkCol, int32_t curIndex) {
234
  if (currentTs > pSliceInfo->win.ekey) {
2,147,483,647✔
235
    return INVALID_TIMESTAMP_REASON_NONE;
1,842,892✔
236
  }
237
  if (pSliceInfo->prevTsSet && currentTs <= pSliceInfo->prevKey.ts) {
2,147,483,647✔
238
    /**
239
      Input data of time slice operator must be ordered by
240
      timestamp ascendingly, except the prev scan.
241
      So prevTs should never be updated to equal or smaller timestamp.
242
    */
243
    return currentTs == pSliceInfo->prevKey.ts ?
302,886,090✔
244
      INVALID_TIMESTAMP_REASON_PREV_TS_EQUAL :
302,886,090✔
245
      INVALID_TIMESTAMP_REASON_PREV_TS_SMALLER;
246
  }
247

248
  SRowKey cur = {.ts = currentTs, .numOfPKs = (pPkCol != NULL) ? 1 : 0};
2,147,483,647✔
249
  if (pPkCol != NULL) {
2,147,483,647✔
250
    cur.pks[0].type = pPkCol->info.type;
2,109,800✔
251
    if (IS_VAR_DATA_TYPE(pPkCol->info.type)) {
2,109,800✔
252
      cur.pks[0].pData = (uint8_t*)colDataGetVarData(pPkCol, curIndex);
663,080✔
253
    } else {
254
      valueSetDatum(cur.pks, pPkCol->info.type,
1,446,720✔
255
                    colDataGetData(pPkCol, curIndex), pPkCol->info.bytes);
1,446,720✔
256
    }
257
  }
258

259
  pSliceInfo->prevTsSet = true;
2,147,483,647✔
260
  tRowKeyAssign(&pSliceInfo->prevKey, &cur);
2,147,483,647✔
261

262
  return INVALID_TIMESTAMP_REASON_NONE;
2,147,483,647✔
263
}
264

265
bool isInterpFunc(SExprInfo* pExprInfo) {
2,147,483,647✔
266
  int32_t functionType = pExprInfo->pExpr->_function.functionType;
2,147,483,647✔
267
  return (functionType == FUNCTION_TYPE_INTERP);
2,147,483,647✔
268
}
269

270
static bool isGroupKeyFunc(SExprInfo* pExprInfo) {
77,351,771✔
271
  int32_t functionType = pExprInfo->pExpr->_function.functionType;
77,351,771✔
272
  return (functionType == FUNCTION_TYPE_GROUP_KEY);
77,351,771✔
273
}
274

275
static bool isSelectGroupConstValueFunc(SExprInfo* pExprInfo) {
743,871✔
276
  int32_t functionType = pExprInfo->pExpr->_function.functionType;
743,871✔
277
  return (functionType == FUNCTION_TYPE_GROUP_CONST_VALUE);
743,871✔
278
}
279

280
bool getIgoreNullRes(SExprSupp* pExprSup) {
16,718,772✔
281
  for (int32_t i = 0; i < pExprSup->numOfExprs; ++i) {
85,469,130✔
282
    SExprInfo* pExprInfo = &pExprSup->pExprInfo[i];
72,495,017✔
283

284
    if (isInterpFunc(pExprInfo)) {
72,495,017✔
285
      for (int32_t j = 0; j < pExprInfo->base.numOfParams; ++j) {
83,274,703✔
286
        SFunctParam* pFuncParam = &pExprInfo->base.pParam[j];
58,047,716✔
287
        if (pFuncParam->type == FUNC_PARAM_TYPE_VALUE) {
58,047,716✔
288
          return pFuncParam->param.i ? true : false;
3,744,126✔
289
        }
290
      }
291
    }
292
  }
293

294
  return false;
12,974,646✔
295
}
296

297
bool checkNullRow(SExprSupp* pExprSup, SSDataBlock* pSrcBlock, int32_t index, bool ignoreNull) {
2,147,483,647✔
298
  if (!ignoreNull) {
2,147,483,647✔
299
    return false;
2,147,483,647✔
300
  }
301

302
  for (int32_t j = 0; j < pExprSup->numOfExprs; ++j) {
2,147,483,647✔
303
    SExprInfo* pExprInfo = &pExprSup->pExprInfo[j];
2,147,483,647✔
304

305
    if (isInterpFunc(pExprInfo)) {
2,147,483,647✔
306
      int32_t          srcSlot = pExprInfo->base.pParam[0].pCol->slotId;
2,147,483,647✔
307
      SColumnInfoData* pSrc = taosArrayGet(pSrcBlock->pDataBlock, srcSlot);
2,147,483,647✔
308

309
      if (colDataIsNull_s(pSrc, index)) {
2,147,483,647✔
310
        return true;
160,039,222✔
311
      }
312
    }
313
  }
314

315
  return false;
2,122,599,080✔
316
}
317

318
static int32_t interpColSetKey(SColumnInfoData* pDst, int32_t rowNum, SGroupKeys* pKey) {
1,323,772,133✔
319
  int32_t code = 0;
1,323,772,133✔
320
  if (pKey->isNull == false) {
1,323,772,133✔
321
    code = colDataSetVal(pDst, rowNum, pKey->pData, false);
1,323,666,199✔
322
  } else {
323
    colDataSetNULL(pDst, rowNum);
106,824✔
324
  }
325
  return code;
1,323,855,105✔
326
}
327

328
/**
329
  @brief Check if the time difference between the fill reference row and
330
  target row exceeds surroundingTime, if so, set the fill row to NULL and use
331
  fillVal to fill.
332
*/
333
static void checkSurroundingTime(const STimeSliceOperatorInfo* pSliceInfo,
513,142,055✔
334
                                 SArray** ppFillRow, SArray* pFillRefRow,
335
                                 int64_t fillRefRowTs) {
336
  *ppFillRow = NULL;
513,142,055✔
337
  uint64_t diff = safe_abs_diff_i64(fillRefRowTs, pSliceInfo->current);
513,142,588✔
338
  if (pSliceInfo->surroundingTime > 0 &&
513,140,466✔
339
      diff > (uint64_t)pSliceInfo->surroundingTime) {
2,985,610✔
340
    return;
962,430✔
341
  }
342
  *ppFillRow = pFillRefRow;
512,178,036✔
343
}
344

345
static bool interpDetermineNearFillRow(STimeSliceOperatorInfo* pSliceInfo, SArray** ppNearRow) {
406,101,178✔
346
  if (!pSliceInfo->isPrevRowSet && !pSliceInfo->isNextRowSet) {
406,101,178✔
347
    *ppNearRow = NULL;
×
348
    return false;
×
349
  }
350
  SGroupKeys *pPrevTsKey = NULL, *pNextTsKey = NULL;
406,102,601✔
351
  int64_t    *pPrevTs = NULL, *pNextTs = NULL;
406,102,601✔
352
  if (pSliceInfo->isPrevRowSet) {
406,102,601✔
353
    pPrevTsKey = taosArrayGet(pSliceInfo->pPrevRow, pSliceInfo->tsCol.slotId);
311,923,560✔
354
    pPrevTs = (int64_t*)pPrevTsKey->pData;
311,921,790✔
355
  }
356
  if (pSliceInfo->isNextRowSet) {
406,098,787✔
357
    pNextTsKey = taosArrayGet(pSliceInfo->pNextRow, pSliceInfo->tsCol.slotId);
383,381,751✔
358
    pNextTs = (int64_t*)pNextTsKey->pData;
383,380,690✔
359
  }
360
  if (!pPrevTsKey) {
406,099,403✔
361
    *ppNearRow = pSliceInfo->pNextRow;
94,179,124✔
362
    checkSurroundingTime(pSliceInfo, ppNearRow, pSliceInfo->pNextRow, *pNextTs);
94,179,124✔
363
  } else if (!pNextTsKey) {
311,920,279✔
364
    *ppNearRow = pSliceInfo->pPrevRow;
22,712,400✔
365
    checkSurroundingTime(pSliceInfo, ppNearRow, pSliceInfo->pPrevRow, *pPrevTs);
22,712,400✔
366
  } else {
367
    if (llabs(pSliceInfo->current - *pPrevTs) <= 
289,207,879✔
368
        llabs(*pNextTs - pSliceInfo->current)) {
289,208,769✔
369
      /* take prev if euqal */
370
      checkSurroundingTime(pSliceInfo, ppNearRow, pSliceInfo->pPrevRow,
161,471,927✔
371
                           *pPrevTs);
372
    } else {
373
      checkSurroundingTime(pSliceInfo, ppNearRow, pSliceInfo->pNextRow,
127,740,387✔
374
                           *pNextTs);
375
    }
376
  }
377
  return true;
406,086,797✔
378
}
379

380
static bool interpDetermineFillRefRow(STimeSliceOperatorInfo* pSliceInfo, SArray** ppOutRow) {
695,675,218✔
381
  bool needFill = false;
695,675,218✔
382
  if (pSliceInfo->fillType == TSDB_FILL_PREV) {
695,675,218✔
383
    if (pSliceInfo->isPrevRowSet) {
67,294,094✔
384
      SGroupKeys* pTsCol = taosArrayGet(pSliceInfo->pPrevRow, pSliceInfo->tsCol.slotId);
67,033,117✔
385
      checkSurroundingTime(pSliceInfo, ppOutRow, pSliceInfo->pPrevRow,
67,033,117✔
386
                           *(int64_t*)pTsCol->pData);
67,033,117✔
387
      needFill = true;
67,033,117✔
388
    }
389
  } else if (pSliceInfo->fillType == TSDB_FILL_NEXT) {
628,381,574✔
390
    if (pSliceInfo->isNextRowSet) {
40,019,315✔
391
      SGroupKeys* pTsCol = taosArrayGet(pSliceInfo->pNextRow, pSliceInfo->tsCol.slotId);
40,019,315✔
392
      checkSurroundingTime(pSliceInfo, ppOutRow, pSliceInfo->pNextRow,
40,019,315✔
393
                           *(int64_t*)pTsCol->pData);
40,019,315✔
394
      needFill = true;
40,019,315✔
395
    }
396
  } else if (pSliceInfo->fillType == TSDB_FILL_NEAR) {
588,361,281✔
397
    needFill = interpDetermineNearFillRow(pSliceInfo, ppOutRow);
406,103,222✔
398
  } else {
399
    needFill = true;
182,258,592✔
400
  }
401
  return needFill;
695,658,798✔
402
}
403

404
static bool genInterpolationResult(STimeSliceOperatorInfo* pSliceInfo, SExprSupp* pExprSup, SSDataBlock* pResBlock,
695,618,184✔
405
                                   SSDataBlock* pSrcBlock, int32_t index, bool beforeTs, SExecTaskInfo* pTaskInfo) {
406
  int32_t code = TSDB_CODE_SUCCESS;
695,618,184✔
407
  int32_t lino = 0;
695,618,184✔
408
  int32_t rows = pResBlock->info.rows;
695,618,184✔
409
  code = timeSliceEnsureBlockCapacity(pSliceInfo, pResBlock);
695,676,558✔
410
  QUERY_CHECK_CODE(code, lino, _end);
695,676,558✔
411
  // todo set the correct primary timestamp column
412

413
  // output the result
414
  int32_t fillColIndex = 0;
695,676,558✔
415
  int32_t groupKeyIndex = 0;
695,676,558✔
416
  bool    hasInterp = true;
695,676,558✔
417
  SArray* pFillRefRow = NULL;
695,676,558✔
418
  bool    needFill = interpDetermineFillRefRow(pSliceInfo, &pFillRefRow);
695,675,668✔
419
  for (int32_t j = 0; j < pExprSup->numOfExprs; ++j) {
2,147,483,647✔
420
    SExprInfo* pExprInfo = &pExprSup->pExprInfo[j];
2,147,483,647✔
421

422
    int32_t          dstSlot = pExprInfo->base.resSchema.slotId;
2,147,483,647✔
423
    SColumnInfoData* pDst = taosArrayGet(pResBlock->pDataBlock, dstSlot);
2,147,483,647✔
424

425
    if (isIrowtsPseudoColumn(pExprInfo)) {
2,147,483,647✔
426
      code = colDataSetVal(pDst, rows, (char*)&pSliceInfo->current, false);
672,662,167✔
427
      QUERY_CHECK_CODE(code, lino, _end);
672,664,128✔
428
      continue;
672,664,128✔
429
    } else if (isIsfilledPseudoColumn(pExprInfo)) {
2,147,483,647✔
430
      bool isFilled = true;
671,836,691✔
431
      code = colDataSetVal(pDst, pResBlock->info.rows, (char*)&isFilled, false);
671,836,691✔
432
      QUERY_CHECK_CODE(code, lino, _end);
671,837,762✔
433
      continue;
671,837,762✔
434
    } else if (!isInterpFunc(pExprInfo) && !isIrowtsOriginPseudoColumn(pExprInfo)) {
1,510,531,091✔
435
      if (isGroupKeyFunc(pExprInfo) || isSelectGroupConstValueFunc(pExprInfo)) {
2,061,563✔
436
        if (pSrcBlock != NULL) {
2,061,563✔
437
          int32_t          srcSlot = pExprInfo->base.pParam[0].pCol->slotId;
1,491,607✔
438
          SColumnInfoData* pSrc = taosArrayGet(pSrcBlock->pDataBlock, srcSlot);
1,491,607✔
439

440
          if (colDataIsNull_s(pSrc, index)) {
2,983,214✔
441
            colDataSetNULL(pDst, pResBlock->info.rows);
6,230✔
442
            continue;
6,230✔
443
          }
444

445
          char* v = colDataGetData(pSrc, index);
1,485,377✔
446
          code = colDataSetVal(pDst, pResBlock->info.rows, v, false);
1,485,377✔
447
          QUERY_CHECK_CODE(code, lino, _end);
1,485,377✔
448
        } else if (!isSelectGroupConstValueFunc(pExprInfo)) {
569,956✔
449
          // use stored group key
450
          SGroupKeys* pkey = taosArrayGet(pSliceInfo->pPrevGroupKeys, groupKeyIndex);
532,681✔
451
          QUERY_CHECK_NULL(pkey, code, lino, _end, terrno);
532,681✔
452
          groupKeyIndex++;
532,681✔
453
          if (pkey->isNull == false) {
532,681✔
454
            code = colDataSetVal(pDst, rows, pkey->pData, false);
532,681✔
455
            QUERY_CHECK_CODE(code, lino, _end);
532,681✔
456
          } else {
457
            colDataSetNULL(pDst, rows);
×
458
          }
459
        } else {
460
          int32_t     srcSlot = pExprInfo->base.pParam[0].pCol->slotId;
37,275✔
461
          SGroupKeys* pkey = taosArrayGet(pSliceInfo->pPrevRow, srcSlot);
37,275✔
462
          if (pkey->isNull == false) {
37,275✔
463
            code = colDataSetVal(pDst, rows, pkey->pData, false);
37,275✔
464
            QUERY_CHECK_CODE(code, lino, _end);
37,275✔
465
          } else {
466
            colDataSetNULL(pDst, rows);
×
467
          }
468
        }
469
      }
470
      continue;
2,055,333✔
471
    }
472

473
    int32_t srcSlot =
1,508,621,803✔
474
        isIrowtsOriginPseudoColumn(pExprInfo) ? pSliceInfo->tsCol.slotId : pExprInfo->base.pParam[0].pCol->slotId;
1,508,518,114✔
475
    switch (pSliceInfo->fillType) {
1,508,621,803✔
476
      case TSDB_FILL_NULL:
39,783,913✔
477
      case TSDB_FILL_NULL_F: {
478
        colDataSetNULL(pDst, rows);
39,783,913✔
479
        break;
39,783,913✔
480
      }
481

482
      case TSDB_FILL_PREV:
1,326,112,078✔
483
      case TSDB_FILL_NEAR:
484
      case TSDB_FILL_NEXT: {
485
        if (!needFill) {
1,326,112,078✔
486
          hasInterp = false;
372,382✔
487
          break;
372,382✔
488
        }
489
        if (pFillRefRow) {
1,325,739,696✔
490
          code = interpColSetKey(pDst, rows, taosArrayGet(pFillRefRow, srcSlot));
1,323,696,962✔
491
          QUERY_CHECK_CODE(code, lino, _end);
1,323,854,127✔
492
          break;
1,323,854,127✔
493
        }
494
        // no fillRefRow, fall through to fill specified values
495
        if (srcSlot == pSliceInfo->tsCol.slotId) {
2,042,734✔
496
          // if is _irowts_origin, there is no value to fill, just set to null
497
          colDataSetNULL(pDst, rows);
659,155✔
498
          break;
659,155✔
499
        }
500
      }
501
      case TSDB_FILL_SET_VALUE:
502
      case TSDB_FILL_SET_VALUE_F: {
503
        SVariant* pVar = &pSliceInfo->pFillColInfo[fillColIndex].fillVal;
107,427,474✔
504

505
        bool isNull = (TSDB_DATA_TYPE_NULL == pVar->nType) ? true : false;
107,486,043✔
506
        if (pDst->info.type == TSDB_DATA_TYPE_FLOAT) {
107,486,043✔
507
          float v = 0;
74,752✔
508
          if (!IS_VAR_DATA_TYPE(pVar->nType)) {
74,752✔
509
            GET_TYPED_DATA(v, float, pVar->nType, &pVar->f, 0);
74,752✔
510
          } else {
511
            v = taosStr2Float(varDataVal(pVar->pz), NULL);
×
512
          }
513
          code = colDataSetVal(pDst, rows, (char*)&v, isNull);
74,752✔
514
          QUERY_CHECK_CODE(code, lino, _end);
74,752✔
515
        } else if (pDst->info.type == TSDB_DATA_TYPE_DOUBLE) {
107,411,291✔
516
          double v = 0;
6,397,591✔
517
          if (!IS_VAR_DATA_TYPE(pVar->nType)) {
6,397,591✔
518
            GET_TYPED_DATA(v, double, pVar->nType, &pVar->d, 0);
6,397,591✔
519
          } else {
520
            v = taosStr2Double(varDataVal(pVar->pz), NULL);
×
521
          }
522
          code = colDataSetVal(pDst, rows, (char*)&v, isNull);
6,397,591✔
523
          QUERY_CHECK_CODE(code, lino, _end);
6,397,591✔
524
        } else if (IS_SIGNED_NUMERIC_TYPE(pDst->info.type)) {
201,758,640✔
525
          int64_t v = 0;
100,744,940✔
526
          if (!IS_VAR_DATA_TYPE(pVar->nType)) {
100,744,940✔
527
            GET_TYPED_DATA(v, int64_t, pVar->nType, &pVar->i, 0);
100,744,940✔
528
          } else {
529
            v = taosStr2Int64(varDataVal(pVar->pz), NULL, 10);
×
530
          }
531
          code = colDataSetVal(pDst, rows, (char*)&v, isNull);
100,744,940✔
532
          QUERY_CHECK_CODE(code, lino, _end);
100,744,940✔
533
        } else if (IS_UNSIGNED_NUMERIC_TYPE(pDst->info.type)) {
339,960✔
534
          uint64_t v = 0;
71,200✔
535
          if (!IS_VAR_DATA_TYPE(pVar->nType)) {
71,200✔
536
            GET_TYPED_DATA(v, uint64_t, pVar->nType, &pVar->u, 0);
71,200✔
537
          } else {
538
            v = taosStr2UInt64(varDataVal(pVar->pz), NULL, 10);
×
539
          }
540
          code = colDataSetVal(pDst, rows, (char*)&v, isNull);
71,200✔
541
          QUERY_CHECK_CODE(code, lino, _end);
71,200✔
542
        } else if (IS_BOOLEAN_TYPE(pDst->info.type)) {
197,560✔
543
          bool v = false;
197,560✔
544
          if (!IS_VAR_DATA_TYPE(pVar->nType)) {
197,560✔
545
            GET_TYPED_DATA(v, bool, pVar->nType, &pVar->i, 0);
197,560✔
546
          } else {
547
            v = taosStr2Int8(varDataVal(pVar->pz), NULL, 10);
×
548
          }
549
          code = colDataSetVal(pDst, rows, (char*)&v, isNull);
197,560✔
550
          QUERY_CHECK_CODE(code, lino, _end);
197,560✔
551
        }
552

553
        ++fillColIndex;
107,486,043✔
554
        break;
107,486,043✔
555
      }
556

557
      case TSDB_FILL_LINEAR: {
36,475,626✔
558
        SFillLinearInfo* pLinearInfo = taosArrayGet(pSliceInfo->pLinearInfo, srcSlot);
36,475,626✔
559

560
        SPoint start = pLinearInfo->start;
36,475,626✔
561
        SPoint end = pLinearInfo->end;
36,475,626✔
562
        SPoint current = {.key = pSliceInfo->current};
36,475,626✔
563

564
        // do not interpolate before ts range, only increate pSliceInfo->current
565
        if (beforeTs && !pLinearInfo->isEndSet) {
36,475,626✔
566
          return true;
227,148✔
567
        }
568

569
        if (!pLinearInfo->isStartSet || !pLinearInfo->isEndSet) {
36,248,478✔
570
          hasInterp = false;
×
571
          break;
×
572
        }
573

574
        if (end.key != INT64_MIN && end.key < pSliceInfo->current) {
36,248,478✔
575
          hasInterp = false;
×
576
          break;
×
577
        }
578

579
        if (start.key == INT64_MIN || end.key == INT64_MIN) {
36,248,478✔
580
          colDataSetNULL(pDst, rows);
6,563,369✔
581
          break;
6,563,369✔
582
        }
583

584
        current.val = taosMemoryCalloc(pLinearInfo->bytes, 1);
29,685,109✔
585
        QUERY_CHECK_NULL(current.val, code, lino, _end, terrno);
29,685,109✔
586
        taosGetLinearInterpolationVal(&current, pLinearInfo->type, &start, &end, pLinearInfo->type,
29,685,109✔
587
                                      typeGetTypeModFromColInfo(&pDst->info));
29,685,109✔
588
        code = colDataSetVal(pDst, rows, (char*)current.val, false);
29,685,109✔
589
        QUERY_CHECK_CODE(code, lino, _end);
29,685,109✔
590

591
        taosMemoryFree(current.val);
29,685,109✔
592
        break;
29,685,109✔
593
      }
594
      case TSDB_FILL_NONE:
×
595
      default:
596
        break;
×
597
    }
598
  }
599

600
  if (hasInterp) {
695,438,740✔
601
    pResBlock->info.rows += 1;
695,177,763✔
602
  }
603

604
_end:
695,444,603✔
605
  if (code != TSDB_CODE_SUCCESS) {
695,443,537✔
UNCOV
606
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
UNCOV
607
    pTaskInfo->code = code;
×
608
    T_LONG_JMP(pTaskInfo->env, code);
×
609
  }
610
  return hasInterp;
695,444,603✔
611
}
612

613
static int32_t addCurrentRowToResult(STimeSliceOperatorInfo* pSliceInfo, SExprSupp* pExprSup, SSDataBlock* pResBlock,
4,094,394✔
614
                                     SSDataBlock* pSrcBlock, int32_t index) {
615
  int32_t code = TSDB_CODE_SUCCESS;
4,094,394✔
616
  int32_t lino = 0;
4,094,394✔
617
  code = timeSliceEnsureBlockCapacity(pSliceInfo, pResBlock);
4,094,394✔
618
  QUERY_CHECK_CODE(code, lino, _end);
4,094,394✔
619
  for (int32_t j = 0; j < pExprSup->numOfExprs; ++j) {
12,608,477✔
620
    SExprInfo* pExprInfo = &pExprSup->pExprInfo[j];
8,514,083✔
621

622
    int32_t          dstSlot = pExprInfo->base.resSchema.slotId;
8,514,083✔
623
    SColumnInfoData* pDst = taosArrayGet(pResBlock->pDataBlock, dstSlot);
8,514,083✔
624

625
    if (isIrowtsPseudoColumn(pExprInfo) || isIrowtsOriginPseudoColumn(pExprInfo)) {
8,514,083✔
626
      code = colDataSetVal(pDst, pResBlock->info.rows, (char*)&pSliceInfo->current, false);
2,250,669✔
627
      QUERY_CHECK_CODE(code, lino, _end);
2,250,669✔
628
    } else if (isIsfilledPseudoColumn(pExprInfo)) {
6,263,414✔
629
      bool isFilled = false;
1,304,568✔
630
      code = colDataSetVal(pDst, pResBlock->info.rows, (char*)&isFilled, false);
1,304,568✔
631
      QUERY_CHECK_CODE(code, lino, _end);
1,304,568✔
632
    } else {
633
      int32_t          srcSlot = pExprInfo->base.pParam[0].pCol->slotId;
4,958,846✔
634
      SColumnInfoData* pSrc = taosArrayGet(pSrcBlock->pDataBlock, srcSlot);
4,958,846✔
635

636
      if (colDataIsNull_s(pSrc, index)) {
9,917,692✔
637
        colDataSetNULL(pDst, pResBlock->info.rows);
1,222,717✔
638
        continue;
1,222,717✔
639
      }
640

641
      char* v = colDataGetData(pSrc, index);
3,736,129✔
642
      code = colDataSetVal(pDst, pResBlock->info.rows, v, false);
3,736,129✔
643
      QUERY_CHECK_CODE(code, lino, _end);
3,736,129✔
644
    }
645
  }
646

647
  pResBlock->info.rows += 1;
4,094,394✔
648

649
_end:
4,094,394✔
650
  if (code != TSDB_CODE_SUCCESS) {
4,094,394✔
651
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
652
  }
653
  return code;
4,094,394✔
654
}
655

656
static int32_t initPrevRowsKeeper(STimeSliceOperatorInfo* pInfo, SSDataBlock* pBlock) {
12,366,093✔
657
  int32_t code = TSDB_CODE_SUCCESS;
12,366,093✔
658
  int32_t lino = 0;
12,366,093✔
659
  if (pInfo->pPrevRow != NULL) {
12,366,093✔
660
    return TSDB_CODE_SUCCESS;
8,691,789✔
661
  }
662

663
  pInfo->pPrevRow = taosArrayInit(4, sizeof(SGroupKeys));
3,674,304✔
664
  if (pInfo->pPrevRow == NULL) {
3,674,304✔
665
    return terrno;
×
666
  }
667

668
  int32_t numOfCols = taosArrayGetSize(pBlock->pDataBlock);
3,674,304✔
669
  for (int32_t i = 0; i < numOfCols; ++i) {
14,275,968✔
670
    SColumnInfoData* pColInfo = taosArrayGet(pBlock->pDataBlock, i);
10,601,664✔
671

672
    SGroupKeys key = {0};
10,601,664✔
673
    key.bytes = pColInfo->info.bytes;
10,601,664✔
674
    key.type = pColInfo->info.type;
10,601,664✔
675
    key.isNull = false;
10,601,664✔
676
    key.pData = taosMemoryCalloc(1, pColInfo->info.bytes);
10,601,664✔
677
    QUERY_CHECK_NULL(key.pData, code, lino, _end, terrno);
10,601,664✔
678
    void* tmp = taosArrayPush(pInfo->pPrevRow, &key);
10,601,664✔
679
    QUERY_CHECK_NULL(tmp, code, lino, _end, terrno);
10,601,664✔
680
  }
681

682
  pInfo->isPrevRowSet = false;
3,674,304✔
683

684
_end:
3,674,304✔
685
  if (code != TSDB_CODE_SUCCESS) {
3,674,304✔
686
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
687
  }
688
  return code;
3,674,304✔
689
}
690

691
static int32_t initNextRowsKeeper(STimeSliceOperatorInfo* pInfo, SSDataBlock* pBlock) {
12,366,093✔
692
  int32_t code = TSDB_CODE_SUCCESS;
12,366,093✔
693
  int32_t lino = 0;
12,366,093✔
694
  if (pInfo->pNextRow != NULL) {
12,366,093✔
695
    return TSDB_CODE_SUCCESS;
8,691,789✔
696
  }
697

698
  pInfo->pNextRow = taosArrayInit(4, sizeof(SGroupKeys));
3,674,304✔
699
  if (pInfo->pNextRow == NULL) {
3,674,304✔
700
    return terrno;
×
701
  }
702

703
  int32_t numOfCols = taosArrayGetSize(pBlock->pDataBlock);
3,674,304✔
704
  for (int32_t i = 0; i < numOfCols; ++i) {
14,275,968✔
705
    SColumnInfoData* pColInfo = taosArrayGet(pBlock->pDataBlock, i);
10,601,664✔
706

707
    SGroupKeys key = {0};
10,601,664✔
708
    key.bytes = pColInfo->info.bytes;
10,601,664✔
709
    key.type = pColInfo->info.type;
10,601,664✔
710
    key.isNull = false;
10,601,664✔
711
    key.pData = taosMemoryCalloc(1, pColInfo->info.bytes);
10,601,664✔
712
    QUERY_CHECK_NULL(key.pData, code, lino, _end, terrno);
10,601,664✔
713

714
    void* tmp = taosArrayPush(pInfo->pNextRow, &key);
10,601,664✔
715
    QUERY_CHECK_NULL(tmp, code, lino, _end, terrno);
10,601,664✔
716
  }
717

718
  pInfo->isNextRowSet = false;
3,674,304✔
719

720
_end:
3,674,304✔
721
  if (code != TSDB_CODE_SUCCESS) {
3,674,304✔
722
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
723
  }
724
  return code;
3,674,304✔
725
}
726

727
static int32_t initFillLinearInfo(STimeSliceOperatorInfo* pInfo, SSDataBlock* pBlock) {
12,366,093✔
728
  int32_t code = TSDB_CODE_SUCCESS;
12,366,093✔
729
  int32_t lino = 0;
12,366,093✔
730
  if (pInfo->pLinearInfo != NULL) {
12,366,093✔
731
    return TSDB_CODE_SUCCESS;
8,691,789✔
732
  }
733

734
  pInfo->pLinearInfo = taosArrayInit(4, sizeof(SFillLinearInfo));
3,674,304✔
735
  if (pInfo->pLinearInfo == NULL) {
3,674,304✔
736
    return terrno;
×
737
  }
738

739
  int32_t numOfCols = taosArrayGetSize(pBlock->pDataBlock);
3,674,304✔
740
  for (int32_t i = 0; i < numOfCols; ++i) {
14,275,968✔
741
    SColumnInfoData* pColInfo = taosArrayGet(pBlock->pDataBlock, i);
10,601,664✔
742

743
    SFillLinearInfo linearInfo = {0};
10,601,664✔
744
    linearInfo.start.key = INT64_MIN;
10,601,664✔
745
    linearInfo.end.key = INT64_MIN;
10,601,664✔
746
    linearInfo.start.val = taosMemoryCalloc(1, pColInfo->info.bytes);
10,601,664✔
747
    QUERY_CHECK_NULL(linearInfo.start.val, code, lino, _end, terrno);
10,601,664✔
748

749
    linearInfo.end.val = taosMemoryCalloc(1, pColInfo->info.bytes);
10,601,664✔
750
    QUERY_CHECK_NULL(linearInfo.end.val, code, lino, _end, terrno);
10,601,664✔
751
    linearInfo.isStartSet = false;
10,601,664✔
752
    linearInfo.isEndSet = false;
10,601,664✔
753
    linearInfo.type = pColInfo->info.type;
10,601,664✔
754
    linearInfo.bytes = pColInfo->info.bytes;
10,601,664✔
755
    void* tmp = taosArrayPush(pInfo->pLinearInfo, &linearInfo);
10,601,664✔
756
    QUERY_CHECK_NULL(tmp, code, lino, _end, terrno);
10,601,664✔
757
  }
758

759
_end:
3,674,304✔
760
  if (code != TSDB_CODE_SUCCESS) {
3,674,304✔
761
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
762
  }
763
  return code;
3,674,304✔
764
}
765

766
static void destroyGroupKey(void* pKey) {
74,136✔
767
  SGroupKeys* key = (SGroupKeys*)pKey;
74,136✔
768
  if (key->pData != NULL) {
74,136✔
769
    taosMemoryFreeClear(key->pData);
74,136✔
770
  }
771
}
74,136✔
772

773
static int32_t initGroupKeyKeeper(STimeSliceOperatorInfo* pInfo, SExprSupp* pExprSup) {
12,366,093✔
774
  if (pInfo->pPrevGroupKeys != NULL) {
12,366,093✔
775
    return TSDB_CODE_SUCCESS;
8,691,789✔
776
  }
777

778
  pInfo->pPrevGroupKeys = taosArrayInit(pExprSup->numOfExprs, sizeof(SGroupKeys));
3,674,304✔
779
  if (pInfo->pPrevGroupKeys == NULL) {
3,674,304✔
780
    return terrno;
×
781
  }
782

783
  for (int32_t i = 0; i < pExprSup->numOfExprs; ++i) {
20,711,588✔
784
    SExprInfo* pExprInfo = &pExprSup->pExprInfo[i];
17,037,284✔
785

786
    if (isGroupKeyFunc(pExprInfo)) {
17,037,284✔
787
      SGroupKeys key = {.bytes = pExprInfo->base.resSchema.bytes,
148,272✔
788
                        .type = pExprInfo->base.resSchema.type,
74,136✔
789
                        .isNull = false,
790
                        .pData = taosMemoryCalloc(1, pExprInfo->base.resSchema.bytes)};
74,136✔
791
      if (!key.pData) {
74,136✔
792
        taosArrayDestroyEx(pInfo->pPrevGroupKeys, destroyGroupKey);
×
793
        pInfo->pPrevGroupKeys = NULL;
×
794
        return terrno;
×
795
      }
796
      if (NULL == taosArrayPush(pInfo->pPrevGroupKeys, &key)) {
148,272✔
797
        taosMemoryFree(key.pData);
×
798
        taosArrayDestroyEx(pInfo->pPrevGroupKeys, destroyGroupKey);
×
799
        pInfo->pPrevGroupKeys = NULL;
×
800
        return terrno;
×
801
      }
802
    }
803
  }
804

805
  return TSDB_CODE_SUCCESS;
3,674,304✔
806
}
807

808
static int32_t initKeeperInfo(STimeSliceOperatorInfo* pInfo, SSDataBlock* pBlock, SExprSupp* pExprSup) {
12,366,093✔
809
  int32_t code;
810
  code = initPrevRowsKeeper(pInfo, pBlock);
12,366,093✔
811
  if (code != TSDB_CODE_SUCCESS) {
12,366,093✔
812
    return TSDB_CODE_FAILED;
×
813
  }
814

815
  code = initNextRowsKeeper(pInfo, pBlock);
12,366,093✔
816
  if (code != TSDB_CODE_SUCCESS) {
12,366,093✔
817
    return TSDB_CODE_FAILED;
×
818
  }
819

820
  code = initFillLinearInfo(pInfo, pBlock);
12,366,093✔
821
  if (code != TSDB_CODE_SUCCESS) {
12,366,093✔
822
    return TSDB_CODE_FAILED;
×
823
  }
824

825
  code = initGroupKeyKeeper(pInfo, pExprSup);
12,366,093✔
826
  if (code != TSDB_CODE_SUCCESS) {
12,366,093✔
827
    return TSDB_CODE_FAILED;
×
828
  }
829

830
  return TSDB_CODE_SUCCESS;
12,366,093✔
831
}
832

833
static void resetPrevRowsKeeper(STimeSliceOperatorInfo* pInfo) {
232,038✔
834
  if (pInfo->pPrevRow == NULL) {
232,038✔
835
    return;
×
836
  }
837

838
  for (int32_t i = 0; i < taosArrayGetSize(pInfo->pLinearInfo); ++i) {
1,000,620✔
839
    SGroupKeys* pKey = taosArrayGet(pInfo->pPrevRow, i);
768,582✔
840
    pKey->isNull = false;
768,582✔
841
  }
842

843
  pInfo->isPrevRowSet = false;
232,038✔
844

845
  return;
232,038✔
846
}
847

848
static void resetNextRowsKeeper(STimeSliceOperatorInfo* pInfo) {
232,038✔
849
  if (pInfo->pNextRow == NULL) {
232,038✔
850
    return;
×
851
  }
852

853
  for (int32_t i = 0; i < taosArrayGetSize(pInfo->pLinearInfo); ++i) {
1,000,620✔
854
    SGroupKeys* pKey = taosArrayGet(pInfo->pPrevRow, i);
768,582✔
855
    pKey->isNull = false;
768,582✔
856
  }
857

858
  pInfo->isNextRowSet = false;
232,038✔
859

860
  return;
232,038✔
861
}
862

863
static void resetFillLinearInfo(STimeSliceOperatorInfo* pInfo) {
232,038✔
864
  if (pInfo->pLinearInfo == NULL) {
232,038✔
865
    return;
×
866
  }
867

868
  for (int32_t i = 0; i < taosArrayGetSize(pInfo->pLinearInfo); ++i) {
1,000,620✔
869
    SFillLinearInfo* pLinearInfo = taosArrayGet(pInfo->pLinearInfo, i);
768,582✔
870
    pLinearInfo->start.key = INT64_MIN;
768,582✔
871
    pLinearInfo->end.key = INT64_MIN;
768,582✔
872
    pLinearInfo->isStartSet = false;
768,582✔
873
    pLinearInfo->isEndSet = false;
768,582✔
874
  }
875

876
  return;
232,038✔
877
}
878

879
static void resetKeeperInfo(STimeSliceOperatorInfo* pInfo) {
232,038✔
880
  resetPrevRowsKeeper(pInfo);
232,038✔
881
  resetNextRowsKeeper(pInfo);
232,038✔
882
  resetFillLinearInfo(pInfo);
232,038✔
883
}
232,038✔
884

885
static bool checkThresholdReached(STimeSliceOperatorInfo* pSliceInfo, int32_t threshold) {
270,799,509✔
886
  SSDataBlock* pResBlock = pSliceInfo->pRes;
270,799,509✔
887
  if (pResBlock->info.rows > threshold) {
270,799,954✔
888
    return true;
19,580✔
889
  }
890

891
  return false;
270,781,435✔
892
}
893

894
static bool checkWindowBoundReached(STimeSliceOperatorInfo* pSliceInfo) {
297,567,736✔
895
  if (pSliceInfo->current > pSliceInfo->win.ekey) {
297,567,736✔
896
    return true;
14,401,244✔
897
  }
898

899
  return false;
283,167,553✔
900
}
901

902
static void saveBlockStatus(STimeSliceOperatorInfo* pSliceInfo, SSDataBlock* pBlock, int32_t curIndex) {
9,790✔
903
  SSDataBlock* pResBlock = pSliceInfo->pRes;
9,790✔
904

905
  SColumnInfoData* pTsCol = taosArrayGet(pBlock->pDataBlock, pSliceInfo->tsCol.slotId);
9,790✔
906
  if (curIndex < pBlock->info.rows - 1) {
9,790✔
907
    pSliceInfo->pRemainRes = pBlock;
9,790✔
908
    pSliceInfo->remainIndex = curIndex + 1;
9,790✔
909
    return;
9,790✔
910
  }
911

912
  // all data in remaining block processed
913
  pSliceInfo->pRemainRes = NULL;
×
914
}
915

916
/**
917
  @brief set the 'get param' for the downstream operator to notify the prev/next
918
  scan when the current timestamp is reached the notifyTs. Here we use the 'get
919
  param' to notify the downstream because the notification is going to impact
920
  the data query flow.
921
  @param pOperator: the current operator(parent operator)
922
  @param notifyTs: the timestamp to notify the downstream operator
923
*/
924
static int32_t setDownstreamOpGetParam(SOperatorInfo* pOperator,
8,755,631✔
925
                                       TSKEY notifyTs) {
926
  int32_t code = TSDB_CODE_SUCCESS;
8,755,631✔
927
  int32_t lino = 0;
8,755,631✔
928

929
  if (pOperator->pDownstreamGetParams == NULL) {
8,755,631✔
930
    pOperator->pDownstreamGetParams =
3,228,958✔
931
      taosMemoryCalloc(pOperator->numOfDownstream, POINTER_BYTES);
3,228,958✔
932
    QUERY_CHECK_NULL(pOperator->pDownstreamGetParams, code, lino, _end,
3,228,958✔
933
                     terrno);
934
  }
935

936
  for (int32_t i = 0; i < pOperator->numOfDownstream; ++i) {
17,511,262✔
937
    SOperatorInfo* pDownstream = pOperator->pDownstream[i];
8,755,631✔
938
    if (pDownstream->operatorType != QUERY_NODE_PHYSICAL_PLAN_TABLE_SCAN &&
8,755,631✔
939
        pDownstream->operatorType != QUERY_NODE_PHYSICAL_PLAN_EXCHANGE) {
5,530,314✔
940
      /**
941
        Only table scan and exchange operator are supported right now.
942
      */
943
      qWarn("%s, %s only table scan and exchange operators are supported "
1,141,834✔
944
             "for notify right now, but got %d, skip notify step done",
945
             GET_TASKID(pOperator->pTaskInfo), __func__,
946
             pDownstream->operatorType);
947
      continue;
1,141,834✔
948
    }
949
    SOperatorParam* pParam = pOperator->pDownstreamGetParams[i];
7,613,797✔
950
    if (pParam == NULL) {
7,613,797✔
951
      pParam = (SOperatorParam*)taosMemoryCalloc(1, sizeof(SOperatorParam));
7,613,797✔
952
      QUERY_CHECK_NULL(pParam, code, lino, _end, terrno);
7,613,797✔
953
    }
954

955
    if (pParam->value == NULL) {
7,613,797✔
956
      void* tsParam = NULL;
7,613,797✔
957
      switch (pDownstream->operatorType) {
7,613,797✔
958
        case QUERY_NODE_PHYSICAL_PLAN_TABLE_SCAN: {
3,225,317✔
959
          tsParam = taosMemoryCalloc(1, sizeof(STableScanOperatorParam));
3,225,317✔
960
          QUERY_CHECK_NULL(tsParam, code, lino, _end, terrno);
3,225,317✔
961
          break;
3,225,317✔
962
        }
963
        case QUERY_NODE_PHYSICAL_PLAN_EXCHANGE: {
4,388,480✔
964
          tsParam = taosMemoryCalloc(1, sizeof(SExchangeOperatorParam));
4,388,480✔
965
          QUERY_CHECK_NULL(tsParam, code, lino, _end, terrno);
4,388,480✔
966
          break;
4,388,480✔
967
        }
968
        default:
×
969
          break;
×
970
      }
971
      pParam->value = tsParam;
7,613,797✔
972
    }
973

974
    switch (pDownstream->operatorType) {
7,613,797✔
975
      case QUERY_NODE_PHYSICAL_PLAN_TABLE_SCAN: {
3,225,317✔
976
        STableScanOperatorParam* p = (STableScanOperatorParam*)pParam->value;
3,225,317✔
977
        p->paramType = NOTIFY_TYPE_SCAN_PARAM;
3,225,317✔
978
        p->notifyToProcess = true;
3,225,317✔
979
        p->notifyTs = notifyTs;
3,225,317✔
980
        break;
3,225,317✔
981
      }
982
      case QUERY_NODE_PHYSICAL_PLAN_EXCHANGE: {
4,388,480✔
983
        SExchangeOperatorParam* p = (SExchangeOperatorParam*)pParam->value;
4,388,480✔
984
        p->multiParams = false;
4,388,480✔
985
        p->basic.paramType = NOTIFY_TYPE_EXCHANGE_PARAM;
4,388,480✔
986
        p->basic.notifyTs = notifyTs;
4,388,480✔
987
        break;
4,388,480✔
988
      }
UNCOV
989
      default: {
×
990
        /**
991
          Only table scan and exchange operator are supported right now.
992
        */
UNCOV
993
        qWarn("%s, %s only table scan and exchange operators are supported "
×
994
               "for notify right now, but got %d, skip notify step done",
995
               GET_TASKID(pOperator->pTaskInfo), __func__,
996
               pDownstream->operatorType);
997
        continue;
×
998
      }
999
    }
1000

1001
    pParam->opType = pDownstream->operatorType;
7,613,797✔
1002
    pParam->downstreamIdx = i;
7,613,797✔
1003
    pParam->pChildren = NULL;
7,613,797✔
1004
    pParam->reUse = true;
7,613,797✔
1005
    pOperator->pDownstreamGetParams[i] = pParam;
7,613,797✔
1006
  }
1007

1008
_end:
8,755,631✔
1009
  if (code != TSDB_CODE_SUCCESS) {
8,755,631✔
1010
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
1011
  }
1012
  return code;
8,755,631✔
1013
}
1014

1015
static int64_t getNextTimestamp(int64_t current, SInterval* pInterval) {
512,922,904✔
1016
  return taosTimeAdd(current, pInterval->interval,
512,923,882✔
1017
                     pInterval->intervalUnit, pInterval->precision, NULL);
512,922,904✔
1018
}
1019

1020
static void doTimesliceImpl(SOperatorInfo* pOperator,
12,366,093✔
1021
                            STimeSliceOperatorInfo* pSliceInfo,
1022
                            SSDataBlock* pBlock, SExecTaskInfo* pTaskInfo,
1023
                            bool ignoreNull) {
1024
  int32_t      code = TSDB_CODE_SUCCESS;
12,366,093✔
1025
  int32_t      lino = 0;
12,366,093✔
1026
  SSDataBlock* pResBlock = pSliceInfo->pRes;
12,366,093✔
1027
  SInterval*   pInterval = &pSliceInfo->interval;
12,366,093✔
1028

1029
  SColumnInfoData* pTsCol = taosArrayGet(pBlock->pDataBlock,
12,366,093✔
1030
                                         pSliceInfo->tsCol.slotId);
12,366,093✔
1031
  SColumnInfoData* pPkCol = NULL;
12,366,093✔
1032

1033
  if (pSliceInfo->hasPk) {
12,366,093✔
1034
    pPkCol = taosArrayGet(pBlock->pDataBlock, pSliceInfo->pkCol.slotId);
105,490✔
1035
  }
1036

1037
  int32_t i = (pSliceInfo->pRemainRes == NULL) ? 0 : pSliceInfo->remainIndex;
12,366,093✔
1038
  for (; i < pBlock->info.rows; ++i) {
2,147,483,647✔
1039
    int64_t ts = *(int64_t*)colDataGetData(pTsCol, i);
2,147,483,647✔
1040

1041
    if (checkNullRow(&pOperator->exprSupp, pBlock, i, ignoreNull)) {
2,147,483,647✔
1042
      continue;
160,038,689✔
1043
    }
1044

1045
    if ((!pSliceInfo->prevNotified && ts < pSliceInfo->win.skey) ||
2,147,483,647✔
1046
        (!pSliceInfo->nextNotified && ts > pSliceInfo->win.ekey)) {
2,147,483,647✔
1047
      code = setDownstreamOpGetParam(pOperator, ts);
4,402,952✔
1048
      QUERY_CHECK_CODE(code, lino, _end);
4,402,952✔
1049
      if (ts < pSliceInfo->win.skey) {
4,402,952✔
1050
        pSliceInfo->prevNotified = true;
2,560,060✔
1051
      } else {
1052
        pSliceInfo->nextNotified = true;
1,842,892✔
1053
      }
1054
    }
1055

1056
    EInvalidTimestampReason invalidReason = isInvalidTimestamp(pSliceInfo, ts,
2,147,483,647✔
1057
                                                               pPkCol, i);
1058
    if (invalidReason != INVALID_TIMESTAMP_REASON_NONE) {
2,147,483,647✔
1059
      if (invalidReason == INVALID_TIMESTAMP_REASON_PREV_TS_EQUAL) {
302,886,090✔
1060
        continue;
297,685,523✔
1061
      } else if (invalidReason == INVALID_TIMESTAMP_REASON_PREV_TS_SMALLER) {
5,200,567✔
1062
        break;
5,200,567✔
1063
      }
1064
    }
1065

1066
    if (ts == pSliceInfo->current) {
2,147,483,647✔
1067
      code = addCurrentRowToResult(pSliceInfo, &pOperator->exprSupp,
896,652✔
1068
                                   pResBlock, pBlock, i);
1069
      QUERY_CHECK_CODE(code, lino, _end);
896,652✔
1070

1071
      doKeepPrevRows(pSliceInfo, pBlock, i);
896,652✔
1072
      doKeepLinearInfo(pSliceInfo, pBlock, i);
896,652✔
1073

1074
      pSliceInfo->current = getNextTimestamp(pSliceInfo->current, pInterval);
896,652✔
1075
      if (checkWindowBoundReached(pSliceInfo)) {
896,652✔
1076
        break;
76,460✔
1077
      }
1078
      if (checkThresholdReached(pSliceInfo, pOperator->resultInfo.threshold)) {
820,192✔
1079
        saveBlockStatus(pSliceInfo, pBlock, i);
×
1080
        return;
×
1081
      }
1082
    } else if (ts < pSliceInfo->current) {
2,147,483,647✔
1083
      doKeepPrevRows(pSliceInfo, pBlock, i);
2,147,483,647✔
1084
      doKeepLinearInfo(pSliceInfo, pBlock, i);
2,147,483,647✔
1085
    } else {  /* ts > pSliceInfo->current */
1086
      doKeepNextRows(pSliceInfo, pBlock, i);
263,236,102✔
1087
      doKeepLinearInfo(pSliceInfo, pBlock, i);
263,233,261✔
1088

1089
      while (pSliceInfo->current < ts &&
772,057,346✔
1090
             pSliceInfo->current <= pSliceInfo->win.ekey) {
509,611,491✔
1091
        if (!genInterpolationResult(pSliceInfo, &pOperator->exprSupp,
508,827,904✔
1092
                                    pResBlock, pBlock, i, true, pTaskInfo) &&
256,972✔
1093
            pSliceInfo->fillType == TSDB_FILL_LINEAR) {
256,972✔
1094
          break;
1095
        } else {
1096
          pSliceInfo->current = getNextTimestamp(pSliceInfo->current,
508,828,510✔
1097
                                                 pInterval);
1098
        }
1099
      }
1100

1101
      // add current row if timestamp matches
1102
      if (ts == pSliceInfo->current &&
263,191,051✔
1103
          pSliceInfo->current <= pSliceInfo->win.ekey) {
3,248,437✔
1104
        code = addCurrentRowToResult(pSliceInfo, &pOperator->exprSupp,
3,197,742✔
1105
                                     pResBlock, pBlock, i);
1106
        QUERY_CHECK_CODE(code, lino, _end);
3,197,742✔
1107

1108
        pSliceInfo->current = getNextTimestamp(pSliceInfo->current, pInterval);
3,197,742✔
1109
      }
1110
      doKeepPrevRows(pSliceInfo, pBlock, i);
263,228,914✔
1111

1112
      if (checkWindowBoundReached(pSliceInfo)) {
263,236,102✔
1113
        break;
2,771,483✔
1114
      }
1115
      if (checkThresholdReached(pSliceInfo, pOperator->resultInfo.threshold)) {
260,463,201✔
1116
        saveBlockStatus(pSliceInfo, pBlock, i);
9,790✔
1117
        return;
9,790✔
1118
      }
1119
    }
1120
  }
1121

1122
  // if reached here, meaning block processing finished naturally,
1123
  // or interpolation reach window upper bound
1124
  pSliceInfo->pRemainRes = NULL;
12,344,461✔
1125

1126
_end:
12,356,303✔
1127
  if (code != TSDB_CODE_SUCCESS) {
12,356,303✔
1128
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
1129
    pTaskInfo->code = code;
×
1130
    T_LONG_JMP(pTaskInfo->env, code);
×
1131
  }
1132
}
1133

1134
static void genInterpAfterDataBlock(STimeSliceOperatorInfo* pSliceInfo, SOperatorInfo* pOperator, int32_t index) {
3,915,587✔
1135
  SSDataBlock* pResBlock = pSliceInfo->pRes;
3,915,587✔
1136
  SInterval*   pInterval = &pSliceInfo->interval;
3,915,587✔
1137

1138
  if (pSliceInfo->fillType == TSDB_FILL_NEXT || pSliceInfo->fillType == TSDB_FILL_LINEAR ||
3,915,587✔
1139
      pSliceInfo->pPrevGroupKeys == NULL) {
3,535,814✔
1140
    return;
389,068✔
1141
  }
1142

1143
  while (pSliceInfo->current <= pSliceInfo->win.ekey) {
190,374,288✔
1144
    (void)genInterpolationResult(pSliceInfo, &pOperator->exprSupp, pResBlock, NULL, index, false, pOperator->pTaskInfo);
186,847,769✔
1145
    pSliceInfo->current =
186,847,769✔
1146
        taosTimeAdd(pSliceInfo->current, pInterval->interval, pInterval->intervalUnit, pInterval->precision, NULL);
186,847,324✔
1147
  }
1148
}
1149

1150
static int32_t copyPrevGroupKey(SExprSupp* pExprSup, SArray* pGroupKeys, SSDataBlock* pSrcBlock) {
12,366,093✔
1151
  int32_t groupKeyIdx = 0;
12,366,093✔
1152
  for (int32_t j = 0; j < pExprSup->numOfExprs; ++j) {
70,619,017✔
1153
    SExprInfo* pExprInfo = &pExprSup->pExprInfo[j];
58,252,924✔
1154

1155
    if (isGroupKeyFunc(pExprInfo)) {
58,252,924✔
1156
      int32_t     srcSlot = pExprInfo->base.pParam[0].pCol->slotId;
1,316,278✔
1157
      SGroupKeys* pGroupKey = taosArrayGet(pGroupKeys, groupKeyIdx);
1,316,278✔
1158
      if (pGroupKey == NULL) {
1,316,278✔
1159
        return terrno;
×
1160
      }
1161
      groupKeyIdx++;
1,316,278✔
1162
      SColumnInfoData* pSrc = taosArrayGet(pSrcBlock->pDataBlock, srcSlot);
1,316,278✔
1163

1164
      if (colDataIsNull_s(pSrc, 0)) {
1,316,278✔
1165
        pGroupKey->isNull = true;
×
1166
        break;
×
1167
      }
1168

1169
      char* v = colDataGetData(pSrc, 0);
1,316,278✔
1170
      if (IS_VAR_DATA_TYPE(pGroupKey->type)) {
1,316,278✔
1171
        if (IS_STR_DATA_BLOB(pGroupKey->type)) {
1,223,796✔
1172
          memcpy(pGroupKey->pData, v, blobDataTLen(v));
×
1173
        } else {
1174
          memcpy(pGroupKey->pData, v, varDataTLen(v));
1,223,796✔
1175
        }
1176
      } else {
1177
        memcpy(pGroupKey->pData, v, pGroupKey->bytes);
92,482✔
1178
      }
1179

1180
      pGroupKey->isNull = false;
1,316,278✔
1181
    }
1182
  }
1183
  return TSDB_CODE_SUCCESS;
12,366,093✔
1184
}
1185

1186
static void resetTimesliceInfo(STimeSliceOperatorInfo* pSliceInfo) {
232,038✔
1187
  pSliceInfo->current = pSliceInfo->win.skey;
232,038✔
1188
  pSliceInfo->prevTsSet = false;
232,038✔
1189
  pSliceInfo->prevNotified = false;
232,038✔
1190
  pSliceInfo->nextNotified = false;
232,038✔
1191
  resetKeeperInfo(pSliceInfo);
232,038✔
1192
}
232,038✔
1193

1194
static void doHandleTimeslice(SOperatorInfo* pOperator, SSDataBlock* pBlock) {
16,718,772✔
1195
  int32_t code = TSDB_CODE_SUCCESS;
16,718,772✔
1196
  int32_t lino = 0;
16,718,772✔
1197
  SExecTaskInfo* pTaskInfo = pOperator->pTaskInfo;
16,718,772✔
1198

1199
  STimeSliceOperatorInfo* pSliceInfo = pOperator->info;
16,718,772✔
1200
  SExprSupp*              pSup = &pOperator->exprSupp;
16,718,772✔
1201
  bool                    ignoreNull = getIgoreNullRes(pSup);
16,718,772✔
1202
  int32_t                 order = TSDB_ORDER_ASC;
16,718,772✔
1203

1204
  if (pSup->pFilterInfo != NULL) {
16,718,772✔
1205
    filterSetExecContext(pSup->pFilterInfo, pTaskInfo, isTaskKilled);
890✔
1206
  }
1207

1208
  if (checkWindowBoundReached(pSliceInfo)) {
16,718,772✔
1209
    code = setDownstreamOpGetParam(pOperator, pSliceInfo->win.ekey + 1);
4,352,679✔
1210
    QUERY_CHECK_CODE(code, lino, _end);
4,352,679✔
1211
    goto _end;
4,352,679✔
1212
  }
1213

1214
  code = initKeeperInfo(pSliceInfo, pBlock, &pOperator->exprSupp);
12,366,093✔
1215
  QUERY_CHECK_CODE(code, lino, _end);
12,366,093✔
1216

1217
  if (pSliceInfo->scalarSup.pExprInfo != NULL) {
12,366,093✔
1218
    SExprSupp* pExprSup = &pSliceInfo->scalarSup;
107,715✔
1219
    code = projectApplyFunctions(pExprSup->pExprInfo, pBlock, pBlock,
215,430✔
1220
                                 pExprSup->pCtx, pExprSup->numOfExprs, NULL,
1221
                                 GET_STM_RTINFO(pOperator->pTaskInfo), pOperator->pTaskInfo);
107,715✔
1222
    QUERY_CHECK_CODE(code, lino, _end);
107,715✔
1223
  }
1224

1225
  // the pDataBlock are always the same one, no need to call this again
1226
  code = setInputDataBlock(pSup, pBlock, order, MAIN_SCAN, true);
12,366,093✔
1227
  QUERY_CHECK_CODE(code, lino, _end);
12,366,093✔
1228
  doTimesliceImpl(pOperator, pSliceInfo, pBlock, pTaskInfo, ignoreNull);
12,366,093✔
1229
  code = copyPrevGroupKey(&pOperator->exprSupp, pSliceInfo->pPrevGroupKeys, pBlock);
12,366,093✔
1230
  QUERY_CHECK_CODE(code, lino, _end);
12,365,565✔
1231

1232
_end:
12,365,565✔
1233
  if (code != TSDB_CODE_SUCCESS) {
16,718,244✔
1234
    qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(code));
×
1235
    T_LONG_JMP(pTaskInfo->env, code);
×
1236
  }
1237
}
16,718,244✔
1238

1239
static int32_t doTimesliceNext(SOperatorInfo* pOperator, SSDataBlock** ppRes) {
7,503,719✔
1240
  int32_t        code = TSDB_CODE_SUCCESS;
7,503,719✔
1241
  int32_t        lino = 0;
7,503,719✔
1242
  SExecTaskInfo* pTaskInfo = pOperator->pTaskInfo;
7,503,719✔
1243
  STimeSliceOperatorInfo* pSliceInfo = pOperator->info;
7,503,719✔
1244
  if (pOperator->status == OP_EXEC_DONE) {
7,503,719✔
1245
    (*ppRes) = NULL;
911,584✔
1246
    return code;
911,584✔
1247
  } else if (pOperator->status == OP_NOT_OPENED && pSliceInfo->pWin) {
6,592,135✔
1248
    code = streamCalcCurrWinTimeRange((STimeRangeNode*)pSliceInfo->pWin,
178,248✔
1249
                                       &pTaskInfo->pStreamRuntimeInfo->funcInfo,
89,124✔
1250
                                       &pSliceInfo->win, NULL, 3);
1251
    QUERY_CHECK_CODE(code, lino, _finished);
89,124✔
1252
    OPTR_SET_OPENED(pOperator);    
89,124✔
1253
    pSliceInfo->current = pSliceInfo->win.skey;
89,124✔
1254
  }
1255

1256
  SSDataBlock* pResBlock = pSliceInfo->pRes;
6,592,135✔
1257
  blockDataCleanup(pResBlock);
6,592,135✔
1258

1259
  if (IS_STREAM_MODE(pTaskInfo)) {
6,592,135✔
1260
    /**
1261
      For stream calculation, the interp operator is triggered by the window,
1262
      so we need to reset the notified status for each window.
1263
    */
1264
    pSliceInfo->prevNotified = false;
181,431✔
1265
    pSliceInfo->nextNotified = false;
181,431✔
1266
  }
1267

1268
  while (1) {
1269
    if (pSliceInfo->pNextGroupRes != NULL) {
6,761,169✔
1270
      doHandleTimeslice(pOperator, pSliceInfo->pNextGroupRes);
230,703✔
1271
      if (checkWindowBoundReached(pSliceInfo) ||
333,654✔
1272
          checkThresholdReached(pSliceInfo, pOperator->resultInfo.threshold)) {
102,951✔
1273
        code = doFilter(pResBlock, pOperator->exprSupp.pFilterInfo, NULL, NULL);
127,752✔
1274
        QUERY_CHECK_CODE(code, lino, _finished);
127,752✔
1275
        if (pSliceInfo->pRemainRes == NULL) {
127,752✔
1276
          pSliceInfo->pNextGroupRes = NULL;
127,752✔
1277
        }
1278
        if (pResBlock->info.rows != 0) {
127,752✔
1279
          goto _finished;
87,360✔
1280
        } else {
1281
          // after fillter if result block has 0 rows, go back to
1282
          // process pNextGroupRes again for unfinished data
1283
          continue;
40,392✔
1284
        }
1285
      }
1286
      pSliceInfo->pNextGroupRes = NULL;
102,951✔
1287
    }
1288

1289
    while (1) {
13,769,711✔
1290
      SSDataBlock* pBlock = pSliceInfo->pRemainRes ?
20,403,128✔
1291
        pSliceInfo->pRemainRes : getNextBlockFromDownstream(pOperator, 0);
20,403,128✔
1292
      if (pBlock == NULL) {
20,403,211✔
1293
        setOperatorCompleted(pOperator);
3,683,104✔
1294
        break;
3,683,549✔
1295
      }
1296
      printDataBlock(pBlock, "doTimesliceNext",
33,440,214✔
1297
                    GET_TASKID(pOperator->pTaskInfo),
16,720,107✔
1298
                    pOperator->pTaskInfo->id.queryId);
16,720,107✔
1299

1300
      pResBlock->info.scanFlag = pBlock->info.scanFlag;
16,720,107✔
1301
      if (pSliceInfo->groupId == 0 && pBlock->info.id.groupId != 0) {
16,720,107✔
1302
        pSliceInfo->groupId = pBlock->info.id.groupId;
83,728✔
1303
      } else {
1304
        if (pSliceInfo->groupId != pBlock->info.id.groupId) {
16,636,379✔
1305
          pSliceInfo->groupId = pBlock->info.id.groupId;
232,038✔
1306
          pSliceInfo->pNextGroupRes = pBlock;
232,038✔
1307
          break;
232,038✔
1308
        }
1309
      }
1310

1311
      doHandleTimeslice(pOperator, pBlock);
16,488,069✔
1312
      if (checkWindowBoundReached(pSliceInfo) ||
25,902,740✔
1313
          checkThresholdReached(pSliceInfo, pOperator->resultInfo.threshold)) {
9,414,671✔
1314
        code = doFilter(pResBlock, pOperator->exprSupp.pFilterInfo, NULL, NULL);
7,083,188✔
1315
        QUERY_CHECK_CODE(code, lino, _finished);
7,082,660✔
1316
        if (pResBlock->info.rows != 0) {
7,082,660✔
1317
          goto _finished;
2,717,830✔
1318
        }
1319
      }
1320
    }
1321
    // post work for a specific group
1322

1323
    // check if need to interpolate after last datablock
1324
    // except for fill(next), fill(linear)
1325
    genInterpAfterDataBlock(pSliceInfo, pOperator, 0);
3,915,587✔
1326

1327
    code = doFilter(pResBlock, pOperator->exprSupp.pFilterInfo, NULL, NULL);
3,915,587✔
1328
    QUERY_CHECK_CODE(code, lino, _finished);
3,915,587✔
1329
    if (pOperator->status == OP_EXEC_DONE) {
3,915,587✔
1330
      break;
3,683,549✔
1331
    }
1332

1333
    // restore initial value for next group
1334
    resetTimesliceInfo(pSliceInfo);
232,038✔
1335
    if (pResBlock->info.rows != 0) {
232,038✔
1336
      break;
103,396✔
1337
    }
1338
  }
1339

1340
_finished:
6,592,135✔
1341
  // restore the value
1342
  setTaskStatus(pOperator->pTaskInfo, TASK_COMPLETED);
6,592,135✔
1343
  if (pResBlock->info.rows == 0) {
6,592,135✔
1344
    pOperator->status = OP_EXEC_DONE;
2,768,850✔
1345
  }
1346
  if (code != TSDB_CODE_SUCCESS) {
6,592,135✔
1347
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
1348
    pTaskInfo->code = code;
×
1349
    T_LONG_JMP(pTaskInfo->env, code);
×
1350
  }
1351

1352
  (*ppRes) = pResBlock->info.rows == 0 ? NULL : pResBlock;
6,592,135✔
1353
  return code;
6,592,135✔
1354
}
1355

1356
static int32_t extractPkColumnFromFuncs(SNodeList* pFuncs, bool* pHasPk,
3,610,219✔
1357
                                        SColumn* pPkColumn) {
1358
  SNode* pNode;
1359
  FOREACH(pNode, pFuncs) {
20,249,656✔
1360
    if ((nodeType(pNode) == QUERY_NODE_TARGET) &&
16,730,385✔
1361
        (nodeType(((STargetNode*)pNode)->pExpr) == QUERY_NODE_FUNCTION)) {
16,730,385✔
1362
      SFunctionNode* pFunc = (SFunctionNode*)((STargetNode*)pNode)->pExpr;
16,730,385✔
1363
      if (fmIsInterpFunc(pFunc->funcId) && pFunc->hasPk) {
16,729,857✔
1364
        SNode* pNode2 = (pFunc->pParameterList->pTail->pNode);
90,420✔
1365
        if ((nodeType(pNode2) == QUERY_NODE_COLUMN) &&
90,420✔
1366
            ((SColumnNode*)pNode2)->isPk) {
90,420✔
1367
          *pHasPk = true;
90,420✔
1368
          *pPkColumn = extractColumnFromColumnNode((SColumnNode*)pNode2);
90,420✔
1369
          break;
90,420✔
1370
        }
1371
      }
1372
    }
1373
  }
1374
  return TSDB_CODE_SUCCESS;
3,610,219✔
1375
}
1376

1377
static int32_t resetTimeSliceOperState(SOperatorInfo* pOper) {
96,551✔
1378
  STimeSliceOperatorInfo* pInfo = pOper->info;
96,551✔
1379
  SExecTaskInfo*           pTaskInfo = pOper->pTaskInfo;
96,551✔
1380
  SInterpFuncPhysiNode* pPhynode = (SInterpFuncPhysiNode*)pOper->pPhyNode;
96,551✔
1381
  pOper->status = OP_NOT_OPENED;
96,018✔
1382

1383
  setTaskStatus(pOper->pTaskInfo, TASK_NOT_COMPLETED);
96,551✔
1384

1385
  int32_t code = resetExprSupp(&pOper->exprSupp, pTaskInfo, pPhynode->pFuncs,
96,018✔
1386
                               NULL, &pTaskInfo->storageAPI.functionStore);
1387
  if (code == 0) {
96,551✔
1388
    code = resetExprSupp(&pInfo->scalarSup, pTaskInfo, pPhynode->pExprs, NULL,
96,551✔
1389
                         &pTaskInfo->storageAPI.functionStore);
1390
  }
1391

1392
  pInfo->current = pInfo->win.skey;
96,551✔
1393
  pInfo->prevTsSet = false;
96,551✔
1394
  pInfo->prevKey.ts = INT64_MIN;
96,551✔
1395
  pInfo->groupId = 0;
96,551✔
1396
  pInfo->pNextGroupRes = NULL;
96,551✔
1397
  pInfo->pRemainRes = NULL;
96,551✔
1398
  pInfo->remainIndex = 0;
96,551✔
1399

1400
  if (pInfo->hasPk) {
96,551✔
1401
    pInfo->prevKey.numOfPKs = 1;
×
1402
    pInfo->prevKey.pks[0].type = pInfo->pkCol.type;
×
1403

1404
    if (IS_VAR_DATA_TYPE(pInfo->pkCol.type)) {
×
1405
      memset(pInfo->prevKey.pks[0].pData, 0, pInfo->pkCol.bytes);
×
1406
    }
1407
  }
1408
  blockDataCleanup(pInfo->pRes);
96,551✔
1409

1410
  for (int32_t i = 0; i < taosArrayGetSize(pInfo->pPrevRow); ++i) {
289,653✔
1411
    SGroupKeys* pKey = taosArrayGet(pInfo->pPrevRow, i);
193,102✔
1412
    taosMemoryFree(pKey->pData);
193,102✔
1413
  }
1414
  taosArrayDestroy(pInfo->pPrevRow);
96,551✔
1415
  pInfo->pPrevRow = NULL;
96,551✔
1416

1417
  for (int32_t i = 0; i < taosArrayGetSize(pInfo->pNextRow); ++i) {
289,653✔
1418
    SGroupKeys* pKey = taosArrayGet(pInfo->pNextRow, i);
193,102✔
1419
    taosMemoryFree(pKey->pData);
193,102✔
1420
  }
1421
  taosArrayDestroy(pInfo->pNextRow);
96,551✔
1422
  pInfo->pNextRow = NULL;
96,551✔
1423

1424
  for (int32_t i = 0; i < taosArrayGetSize(pInfo->pLinearInfo); ++i) {
289,653✔
1425
    SFillLinearInfo* pKey = taosArrayGet(pInfo->pLinearInfo, i);
193,102✔
1426
    taosMemoryFree(pKey->start.val);
193,102✔
1427
    taosMemoryFree(pKey->end.val);
193,102✔
1428
  }
1429
  taosArrayDestroy(pInfo->pLinearInfo);
96,551✔
1430
  pInfo->pLinearInfo = NULL;
96,551✔
1431

1432
  if (pInfo->pPrevGroupKeys) {
96,551✔
1433
    taosArrayDestroyEx(pInfo->pPrevGroupKeys, destroyGroupKey);
96,551✔
1434
    pInfo->pPrevGroupKeys = NULL;
96,551✔
1435
  }
1436

1437
  return code;
96,551✔
1438
}
1439

1440
int32_t createTimeSliceOperatorInfo(SOperatorInfo* downstream,
3,610,219✔
1441
                                    SPhysiNode* pPhyNode,
1442
                                    SExecTaskInfo* pTaskInfo,
1443
                                    SOperatorInfo** pOptrInfo) {
1444
  QRY_PARAM_CHECK(pOptrInfo);
3,610,219✔
1445

1446
  int32_t                 code = 0;
3,610,219✔
1447
  int32_t                 lino = 0;
3,610,219✔
1448
  STimeSliceOperatorInfo* pInfo =
7,220,438✔
1449
    taosMemoryCalloc(1, sizeof(STimeSliceOperatorInfo));
3,610,219✔
1450
  SOperatorInfo* pOperator = taosMemoryCalloc(1, sizeof(SOperatorInfo));
3,610,219✔
1451

1452
  if (pOperator == NULL || pInfo == NULL) {
3,610,219✔
1453
    code = terrno;
×
1454
    goto _error;
×
1455
  }
1456
  initOperatorCostInfo(pOperator);
3,610,219✔
1457

1458
  pOperator->pPhyNode = pPhyNode;
3,610,219✔
1459
  SInterpFuncPhysiNode* pInterpPhyNode = (SInterpFuncPhysiNode*)pPhyNode;
3,610,219✔
1460
  SExprSupp*            pSup = &pOperator->exprSupp;
3,610,219✔
1461

1462
  int32_t    numOfExprs = 0;
3,610,219✔
1463
  SExprInfo* pExprInfo = NULL;
3,610,219✔
1464
  code = createExprInfo(pInterpPhyNode->pFuncs, NULL, &pExprInfo, &numOfExprs);
3,610,219✔
1465
  QUERY_CHECK_CODE(code, lino, _error);
3,610,219✔
1466

1467
  code = initExprSupp(pSup, pExprInfo, numOfExprs,
3,610,219✔
1468
                      &pTaskInfo->storageAPI.functionStore);
1469
  QUERY_CHECK_CODE(code, lino, _error);
3,610,219✔
1470

1471
  if (pInterpPhyNode->pExprs != NULL) {
3,610,219✔
1472
    int32_t    num = 0;
92,645✔
1473
    SExprInfo* pScalarExprInfo = NULL;
92,645✔
1474
    code = createExprInfo(pInterpPhyNode->pExprs, NULL, &pScalarExprInfo, &num);
92,645✔
1475
    QUERY_CHECK_CODE(code, lino, _error);
92,645✔
1476

1477
    code = initExprSupp(&pInfo->scalarSup, pScalarExprInfo, num,
92,645✔
1478
                        &pTaskInfo->storageAPI.functionStore);
1479
    QUERY_CHECK_CODE(code, lino, _error);
92,645✔
1480
  }
1481

1482
  code = filterInitFromNode((SNode*)pInterpPhyNode->node.pConditions,
3,610,219✔
1483
                            &pOperator->exprSupp.pFilterInfo, 0,
1484
                            pTaskInfo->pStreamRuntimeInfo);
3,610,219✔
1485
  QUERY_CHECK_CODE(code, lino, _error);
3,610,219✔
1486

1487
  pInfo->tsCol =
1488
    extractColumnFromColumnNode((SColumnNode*)pInterpPhyNode->pTimeSeries);
3,610,219✔
1489
  code = extractPkColumnFromFuncs(pInterpPhyNode->pFuncs, &pInfo->hasPk,
3,610,219✔
1490
                                  &pInfo->pkCol);
1491
  QUERY_CHECK_CODE(code, lino, _error);
3,610,219✔
1492

1493
  pInfo->fillType = convertFillType(pInterpPhyNode->fillMode);
3,610,219✔
1494
  initResultSizeInfo(&pOperator->resultInfo, 4096);
3,610,219✔
1495

1496
  pInfo->pFillColInfo =
3,610,219✔
1497
    createFillColInfo(pExprInfo, numOfExprs, NULL, 0, NULL, 0,
3,610,219✔
1498
                      (SNodeListNode*)pInterpPhyNode->pFillValues);
3,609,691✔
1499
  QUERY_CHECK_NULL(pInfo->pFillColInfo, code, lino, _error, terrno);
3,610,219✔
1500

1501
  pInfo->pLinearInfo = NULL;
3,610,219✔
1502
  pInfo->pRes = createDataBlockFromDescNode(pPhyNode->pOutputDataBlockDesc);
3,610,219✔
1503
  QUERY_CHECK_NULL(pInfo->pRes, code, lino, _error, terrno);
3,610,219✔
1504
  pInfo->win = pInterpPhyNode->timeRange;
3,610,219✔
1505
  code = nodesCloneNode(pInterpPhyNode->pTimeRange, &pInfo->pWin);
3,610,219✔
1506
  QUERY_CHECK_CODE(code, lino, _error);
3,610,219✔
1507
  pInfo->interval.interval = pInterpPhyNode->interval;
3,610,219✔
1508
  pInfo->current = pInfo->win.skey;
3,610,219✔
1509
  pInfo->prevTsSet = false;
3,610,219✔
1510
  pInfo->prevKey.ts = INT64_MIN;
3,609,691✔
1511
  pInfo->groupId = 0;
3,610,219✔
1512
  pInfo->pPrevGroupKeys = NULL;
3,610,219✔
1513
  pInfo->pNextGroupRes = NULL;
3,609,774✔
1514
  pInfo->pRemainRes = NULL;
3,609,246✔
1515
  pInfo->remainIndex = 0;
3,610,219✔
1516
  pInfo->surroundingTime = pInterpPhyNode->surroundingTime;
3,610,219✔
1517

1518
  if (pInfo->hasPk) {
3,609,246✔
1519
    pInfo->prevKey.numOfPKs = 1;
90,420✔
1520
    pInfo->prevKey.ts = INT64_MIN;
90,420✔
1521
    pInfo->prevKey.pks[0].type = pInfo->pkCol.type;
90,420✔
1522

1523
    if (IS_VAR_DATA_TYPE(pInfo->pkCol.type)) {
90,420✔
1524
      pInfo->prevKey.pks[0].pData = taosMemoryCalloc(1, pInfo->pkCol.bytes);
30,140✔
1525
      QUERY_CHECK_NULL(pInfo->prevKey.pks[0].pData, code, lino, _error, terrno);
30,140✔
1526
    }
1527
  }
1528

1529
  setOperatorInfo(pOperator, "TimeSliceOperator",
3,609,774✔
1530
                  QUERY_NODE_PHYSICAL_PLAN_INTERP_FUNC, false, OP_NOT_OPENED,
1531
                  pInfo, pTaskInfo);
1532
  pOperator->fpSet = createOperatorFpSet(optrDummyOpenFn, doTimesliceNext,
3,610,219✔
1533
                                         NULL, destroyTimeSliceOperatorInfo,
1534
                                         optrDefaultBufFn, NULL,
1535
                                         optrDefaultGetNextExtFn, NULL);
1536

1537
  code = blockDataEnsureCapacity(pInfo->pRes, pOperator->resultInfo.capacity);
3,609,691✔
1538
  QUERY_CHECK_CODE(code, lino, _error);
3,610,219✔
1539

1540
  //  int32_t code = initKeeperInfo(pSliceInfo, pBlock, &pOperator->exprSupp);
1541
  setOperatorResetStateFn(pOperator, resetTimeSliceOperState);
3,610,219✔
1542
  
1543
  code = appendDownstream(pOperator, &downstream, 1);
3,610,219✔
1544
  QUERY_CHECK_CODE(code, lino, _error);
3,610,219✔
1545

1546
  *pOptrInfo = pOperator;
3,610,219✔
1547
  return TSDB_CODE_SUCCESS;
3,610,219✔
1548

1549
_error:
×
1550
  if (code != TSDB_CODE_SUCCESS) {
×
1551
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
1552
  }
1553
  if (pInfo != NULL) destroyTimeSliceOperatorInfo(pInfo);
×
1554
  destroyOperatorAndDownstreams(pOperator, &downstream, 1);
×
1555
  pTaskInfo->code = code;
×
1556
  return code;
×
1557
}
1558

1559
void destroyTimeSliceOperatorInfo(void* param) {
3,610,219✔
1560
  STimeSliceOperatorInfo* pInfo = (STimeSliceOperatorInfo*)param;
3,610,219✔
1561

1562
  blockDataDestroy(pInfo->pRes);
3,610,219✔
1563
  pInfo->pRes = NULL;
3,610,219✔
1564

1565
  for (int32_t i = 0; i < taosArrayGetSize(pInfo->pPrevRow); ++i) {
14,018,781✔
1566
    SGroupKeys* pKey = taosArrayGet(pInfo->pPrevRow, i);
10,408,562✔
1567
    taosMemoryFree(pKey->pData);
10,408,562✔
1568
  }
1569
  taosArrayDestroy(pInfo->pPrevRow);
3,610,219✔
1570

1571
  for (int32_t i = 0; i < taosArrayGetSize(pInfo->pNextRow); ++i) {
14,018,781✔
1572
    SGroupKeys* pKey = taosArrayGet(pInfo->pNextRow, i);
10,408,562✔
1573
    taosMemoryFree(pKey->pData);
10,408,562✔
1574
  }
1575
  taosArrayDestroy(pInfo->pNextRow);
3,610,219✔
1576

1577
  for (int32_t i = 0; i < taosArrayGetSize(pInfo->pLinearInfo); ++i) {
14,018,781✔
1578
    SFillLinearInfo* pKey = taosArrayGet(pInfo->pLinearInfo, i);
10,408,562✔
1579
    taosMemoryFree(pKey->start.val);
10,408,562✔
1580
    taosMemoryFree(pKey->end.val);
10,408,562✔
1581
  }
1582
  taosArrayDestroy(pInfo->pLinearInfo);
3,610,219✔
1583

1584
  if (pInfo->pPrevGroupKeys) {
3,610,219✔
1585
    taosArrayDestroyEx(pInfo->pPrevGroupKeys, destroyGroupKey);
3,577,753✔
1586
    pInfo->pPrevGroupKeys = NULL;
3,577,753✔
1587
  }
1588
  if (pInfo->hasPk && IS_VAR_DATA_TYPE(pInfo->pkCol.type)) {
3,610,219✔
1589
    taosMemoryFreeClear(pInfo->prevKey.pks[0].pData);
30,140✔
1590
  }
1591

1592
  cleanupExprSupp(&pInfo->scalarSup);
3,610,219✔
1593
  if (pInfo->pFillColInfo != NULL) {
3,610,219✔
1594
    for (int32_t i = 0; i < pInfo->pFillColInfo->numOfFillExpr; ++i) {
20,406,892✔
1595
      taosVariantDestroy(&pInfo->pFillColInfo[i].fillVal);
16,796,673✔
1596
    }
1597
    taosMemoryFree(pInfo->pFillColInfo);
3,610,219✔
1598
  }
1599
  nodesDestroyNode(pInfo->pWin);
3,610,219✔
1600
  taosMemoryFreeClear(param);
3,610,219✔
1601
}
3,610,219✔
1602

1603
STrueForInfo* getTrueForInfo(struct SOperatorInfo* pOperator) {
128,104,270✔
1604
  if (pOperator == NULL) {
128,104,270✔
1605
    return NULL;
×
1606
  }
1607

1608
  switch (pOperator->operatorType) {
128,104,270✔
1609
    case QUERY_NODE_PHYSICAL_PLAN_MERGE_STATE:
3,544,755✔
1610
      return &((SStateWindowOperatorInfo*)pOperator->info)->trueForInfo;
3,544,755✔
1611
    case QUERY_NODE_PHYSICAL_PLAN_MERGE_EVENT:
3,682,611✔
1612
      return &((SEventWindowOperatorInfo*)pOperator->info)->trueForInfo;
3,682,611✔
1613
    default:
120,880,855✔
1614
      return NULL;
120,880,855✔
1615
  }
1616
}
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