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

taosdata / TDengine / #3843

08 Apr 2025 10:23AM UTC coverage: 63.077% (+0.4%) from 62.696%
#3843

push

travis-ci

web-flow
fix: clear cache when meta abort (#30674)

155571 of 315083 branches covered (49.37%)

Branch coverage included in aggregate %.

241876 of 315013 relevant lines covered (76.78%)

19243431.01 hits per line

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

74.34
/source/libs/executor/src/streamfilloperator.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 "filter.h"
17
#include "os.h"
18
#include "query.h"
19
#include "taosdef.h"
20
#include "tmsg.h"
21
#include "ttypes.h"
22

23
#include "executorInt.h"
24
#include "streamexecutorInt.h"
25
#include "streamsession.h"
26
#include "streaminterval.h"
27
#include "tcommon.h"
28
#include "thash.h"
29
#include "ttime.h"
30

31
#include "function.h"
32
#include "operator.h"
33
#include "querynodes.h"
34
#include "querytask.h"
35
#include "tdatablock.h"
36
#include "tfill.h"
37

38
#define FILL_POS_INVALID 0
39
#define FILL_POS_START   1
40
#define FILL_POS_MID     2
41
#define FILL_POS_END     3
42

43
TSKEY getNextWindowTs(TSKEY ts, SInterval* pInterval) {
55✔
44
  STimeWindow win = {.skey = ts, .ekey = ts};
55✔
45
  getNextTimeWindow(pInterval, &win, TSDB_ORDER_ASC);
55✔
46
  return win.skey;
55✔
47
}
48

49
TSKEY getPrevWindowTs(TSKEY ts, SInterval* pInterval) {
30✔
50
  STimeWindow win = {.skey = ts, .ekey = ts};
30✔
51
  getNextTimeWindow(pInterval, &win, TSDB_ORDER_DESC);
30✔
52
  return win.skey;
30✔
53
}
54

55
int32_t setRowCell(SColumnInfoData* pCol, int32_t rowId, const SResultCellData* pCell) {
411,325✔
56
  return colDataSetVal(pCol, rowId, pCell->pData, pCell->isNull);
411,325✔
57
}
58

59
SResultCellData* getResultCell(SResultRowData* pRaw, int32_t index) {
163,255✔
60
  if (!pRaw || !pRaw->pRowVal) {
163,255!
61
    return NULL;
×
62
  }
63
  char*            pData = (char*)pRaw->pRowVal;
163,275✔
64
  SResultCellData* pCell = pRaw->pRowVal;
163,275✔
65
  for (int32_t i = 0; i < index; i++) {
1,129,052✔
66
    pData += (pCell->bytes + sizeof(SResultCellData));
965,777✔
67
    pCell = (SResultCellData*)pData;
965,777✔
68
  }
69
  return pCell;
163,275✔
70
}
71

72
void* destroyFillColumnInfo(SFillColInfo* pFillCol, int32_t start, int32_t end) {
693✔
73
  for (int32_t i = start; i < end; i++) {
1,185✔
74
    destroyExprInfo(pFillCol[i].pExpr, 1);
492✔
75
    taosVariantDestroy(&pFillCol[i].fillVal);
492✔
76
  }
77
  if (start < end) {
693✔
78
    taosMemoryFreeClear(pFillCol[start].pExpr);
451!
79
  }
80
  taosMemoryFree(pFillCol);
693!
81
  return NULL;
693✔
82
}
83

84
void destroyStreamFillSupporter(SStreamFillSupporter* pFillSup) {
693✔
85
  if (pFillSup == NULL) {
693!
86
    return;
×
87
  }
88
  pFillSup->pAllColInfo = destroyFillColumnInfo(pFillSup->pAllColInfo, pFillSup->numOfFillCols, pFillSup->numOfAllCols);
693✔
89
  tSimpleHashCleanup(pFillSup->pResMap);
693✔
90
  pFillSup->pResMap = NULL;
693✔
91
  cleanupExprSupp(&pFillSup->notFillExprSup);
693✔
92
  if (pFillSup->cur.pRowVal != pFillSup->prev.pRowVal && pFillSup->cur.pRowVal != pFillSup->next.pRowVal) {
693✔
93
    taosMemoryFree(pFillSup->cur.pRowVal);
169!
94
  }
95
  taosMemoryFree(pFillSup->prev.pRowVal);
693!
96
  taosMemoryFree(pFillSup->next.pRowVal);
693!
97
  taosMemoryFree(pFillSup->nextNext.pRowVal);
693!
98

99
  taosMemoryFree(pFillSup->pOffsetInfo);
693!
100
  taosArrayDestroy(pFillSup->pResultRange);
693✔
101
  pFillSup->pResultRange = NULL;
693✔
102

103
  taosMemoryFree(pFillSup);
693!
104
}
105

106
void destroySPoint(void* ptr) {
103,862✔
107
  SPoint* point = (SPoint*)ptr;
103,862✔
108
  taosMemoryFreeClear(point->val);
103,862!
109
}
103,949✔
110

111
void destroyStreamFillLinearInfo(SStreamFillLinearInfo* pFillLinear) {
693✔
112
  taosArrayDestroyEx(pFillLinear->pEndPoints, destroySPoint);
693✔
113
  taosArrayDestroyEx(pFillLinear->pNextEndPoints, destroySPoint);
693✔
114
  taosMemoryFree(pFillLinear);
693!
115
}
693✔
116

117
void destroyStreamFillInfo(SStreamFillInfo* pFillInfo) {
693✔
118
  if (pFillInfo == NULL) {
693!
119
    return;
×
120
  } 
121
  if (pFillInfo->type == TSDB_FILL_SET_VALUE || pFillInfo->type == TSDB_FILL_SET_VALUE_F ||
693✔
122
      pFillInfo->type == TSDB_FILL_NULL || pFillInfo->type == TSDB_FILL_NULL_F) {
554✔
123
    taosMemoryFreeClear(pFillInfo->pResRow->pRowVal);
305!
124
    taosMemoryFreeClear(pFillInfo->pResRow);
305!
125
    taosMemoryFreeClear(pFillInfo->pNonFillRow->pRowVal);
305!
126
    taosMemoryFreeClear(pFillInfo->pNonFillRow);
305!
127
  }
128
  destroyStreamFillLinearInfo(pFillInfo->pLinearInfo);
693✔
129
  pFillInfo->pLinearInfo = NULL;
693✔
130

131
  taosArrayDestroy(pFillInfo->delRanges);
693✔
132
  taosMemoryFreeClear(pFillInfo->pTempBuff);
693!
133
  taosMemoryFree(pFillInfo);
693!
134
}
135

136
void clearGroupResArray(SGroupResInfo* pGroupResInfo) {
693✔
137
  pGroupResInfo->freeItem = false;
693✔
138
  taosArrayDestroy(pGroupResInfo->pRows);
693✔
139
  pGroupResInfo->pRows = NULL;
693✔
140
  pGroupResInfo->index = 0;
693✔
141
}
693✔
142

143
void destroyStreamFillOperatorInfo(void* param) {
451✔
144
  SStreamFillOperatorInfo* pInfo = (SStreamFillOperatorInfo*)param;
451✔
145
  destroyStreamFillInfo(pInfo->pFillInfo);
451✔
146
  destroyStreamFillSupporter(pInfo->pFillSup);
451✔
147
  blockDataDestroy(pInfo->pRes);
451✔
148
  pInfo->pRes = NULL;
451✔
149
  blockDataDestroy(pInfo->pSrcBlock);
451✔
150
  pInfo->pSrcBlock = NULL;
451✔
151
  blockDataDestroy(pInfo->pDelRes);
451✔
152
  pInfo->pDelRes = NULL;
451✔
153
  taosArrayDestroy(pInfo->matchInfo.pList);
451✔
154
  pInfo->matchInfo.pList = NULL;
451✔
155
  taosArrayDestroy(pInfo->pUpdated);
451✔
156
  clearGroupResArray(&pInfo->groupResInfo);
451✔
157
  taosArrayDestroy(pInfo->pCloseTs);
451✔
158

159
  if (pInfo->stateStore.streamFileStateDestroy != NULL) {
451✔
160
    pInfo->stateStore.streamFileStateDestroy(pInfo->pState->pFileState);
44✔
161
  }
162

163
  if (pInfo->pState != NULL) {
451✔
164
    taosMemoryFreeClear(pInfo->pState);
44!
165
  }
166
  destroyStreamBasicInfo(&pInfo->basic);
451✔
167
  destroyNonBlockAggSupptor(&pInfo->nbSup);
451✔
168

169
  taosMemoryFree(pInfo);
451!
170
}
451✔
171

172
static void resetFillWindow(SResultRowData* pRowData) {
8,002✔
173
  pRowData->key = INT64_MIN;
8,002✔
174
  taosMemoryFreeClear(pRowData->pRowVal);
8,002!
175
}
8,003✔
176

177
static void resetPrevAndNextWindow(SStreamFillSupporter* pFillSup) {
1,815✔
178
  if (pFillSup->cur.pRowVal != pFillSup->prev.pRowVal && pFillSup->cur.pRowVal != pFillSup->next.pRowVal) {
1,815✔
179
    resetFillWindow(&pFillSup->cur);
1,364✔
180
  } else {
181
    pFillSup->cur.key = INT64_MIN;
451✔
182
    pFillSup->cur.pRowVal = NULL;
451✔
183
  }
184
  resetFillWindow(&pFillSup->prev);
1,815✔
185
  resetFillWindow(&pFillSup->next);
1,815✔
186
  resetFillWindow(&pFillSup->nextNext);
1,815✔
187
}
1,816✔
188

189
void getWindowFromDiscBuf(SOperatorInfo* pOperator, TSKEY ts, uint64_t groupId, SStreamFillSupporter* pFillSup) {
1,815✔
190
  SStorageAPI* pAPI = &pOperator->pTaskInfo->storageAPI;
1,815✔
191
  void*        pState = pOperator->pTaskInfo->streamInfo.pState;
1,815✔
192
  resetPrevAndNextWindow(pFillSup);
1,815✔
193

194
  SWinKey key = {.ts = ts, .groupId = groupId};
1,816✔
195
  void*   curVal = NULL;
1,816✔
196
  int32_t curVLen = 0;
1,816✔
197
  bool    hasCurKey = true;
1,816✔
198
  int32_t code = pAPI->stateStore.streamStateFillGet(pState, &key, (void**)&curVal, &curVLen, NULL);
1,816✔
199
  if (code == TSDB_CODE_SUCCESS) {
1,816✔
200
    pFillSup->cur.key = key.ts;
1,761✔
201
    pFillSup->cur.pRowVal = curVal;
1,761✔
202
  } else {
203
    qDebug("streamStateFillGet key failed, Data may be deleted. ts:%" PRId64 ", groupId:%" PRId64, ts, groupId);
55✔
204
    pFillSup->cur.key = ts;
55✔
205
    pFillSup->cur.pRowVal = NULL;
55✔
206
    hasCurKey = false;
55✔
207
  }
208

209
  SStreamStateCur* pCur = pAPI->stateStore.streamStateFillSeekKeyPrev(pState, &key);
1,816✔
210
  SWinKey          preKey = {.ts = INT64_MIN, .groupId = groupId};
1,816✔
211
  void*            preVal = NULL;
1,816✔
212
  int32_t          preVLen = 0;
1,816✔
213
  code = pAPI->stateStore.streamStateFillGetGroupKVByCur(pCur, &preKey, (const void**)&preVal, &preVLen);
1,816✔
214

215
  if (code == TSDB_CODE_SUCCESS) {
1,816✔
216
    pFillSup->prev.key = preKey.ts;
1,426✔
217
    pFillSup->prev.pRowVal = preVal;
1,426✔
218

219
    if (hasCurKey) {
1,426✔
220
      pAPI->stateStore.streamStateCurNext(pState, pCur);
1,371✔
221
    }
222

223
    pAPI->stateStore.streamStateCurNext(pState, pCur);
1,426✔
224
  } else {
225
    pAPI->stateStore.streamStateFreeCur(pCur);
390✔
226
    pCur = pAPI->stateStore.streamStateFillSeekKeyNext(pState, &key);
390✔
227
  }
228

229
  SWinKey nextKey = {.ts = INT64_MIN, .groupId = groupId};
1,816✔
230
  void*   nextVal = NULL;
1,816✔
231
  int32_t nextVLen = 0;
1,816✔
232
  code = pAPI->stateStore.streamStateFillGetGroupKVByCur(pCur, &nextKey, (const void**)&nextVal, &nextVLen);
1,816✔
233
  if (code == TSDB_CODE_SUCCESS) {
1,815✔
234
    pFillSup->next.key = nextKey.ts;
1,041✔
235
    pFillSup->next.pRowVal = nextVal;
1,041✔
236
    if (pFillSup->type == TSDB_FILL_PREV || pFillSup->type == TSDB_FILL_NEXT) {
1,041✔
237
      pAPI->stateStore.streamStateCurNext(pState, pCur);
394✔
238
      SWinKey nextNextKey = {.groupId = groupId};
394✔
239
      void*   nextNextVal = NULL;
394✔
240
      int32_t nextNextVLen = 0;
394✔
241
      code = pAPI->stateStore.streamStateFillGetGroupKVByCur(pCur, &nextNextKey, (const void**)&nextNextVal, &nextNextVLen);
394✔
242
      if (code == TSDB_CODE_SUCCESS) {
394✔
243
        pFillSup->nextNext.key = nextNextKey.ts;
113✔
244
        pFillSup->nextNext.pRowVal = nextNextVal;
113✔
245
      }
246
    }
247
  }
248
  pAPI->stateStore.streamStateFreeCur(pCur);
1,815✔
249
}
1,816✔
250

251
bool hasCurWindow(SStreamFillSupporter* pFillSup) { return pFillSup->cur.key != INT64_MIN; }
×
252
bool hasPrevWindow(SStreamFillSupporter* pFillSup) { return pFillSup->prev.key != INT64_MIN; }
8,607✔
253
bool hasNextWindow(SStreamFillSupporter* pFillSup) { return pFillSup->next.key != INT64_MIN; }
6,166✔
254
static bool hasNextNextWindow(SStreamFillSupporter* pFillSup) { return pFillSup->nextNext.key != INT64_MIN; }
149✔
255

256
static void transBlockToResultRow(const SSDataBlock* pBlock, int32_t rowId, TSKEY ts, SResultRowData* pRowVal) {
1,802✔
257
  int32_t numOfCols = taosArrayGetSize(pBlock->pDataBlock);
1,802✔
258
  for (int32_t i = 0; i < numOfCols; ++i) {
29,408✔
259
    SColumnInfoData* pColData = taosArrayGet(pBlock->pDataBlock, i);
27,609✔
260
    SResultCellData* pCell = getResultCell(pRowVal, i);
27,608✔
261
    if (!colDataIsNull_s(pColData, rowId)) {
55,212!
262
      pCell->isNull = false;
27,606✔
263
      pCell->type = pColData->info.type;
27,606✔
264
      pCell->bytes = pColData->info.bytes;
27,606✔
265
      char* val = colDataGetData(pColData, rowId);
27,606!
266
      if (IS_VAR_DATA_TYPE(pCell->type)) {
27,606!
267
        memcpy(pCell->pData, val, varDataTLen(val));
38✔
268
      } else {
269
        memcpy(pCell->pData, val, pCell->bytes);
27,568✔
270
      }
271
    } else {
272
      pCell->isNull = true;
×
273
    }
274
  }
275
  pRowVal->key = ts;
1,799✔
276
}
1,799✔
277

278
static void calcRowDeltaData(SResultRowData* pEndRow, SArray* pEndPoins, SFillColInfo* pFillCol, int32_t numOfCol) {
294✔
279
  for (int32_t i = 0; i < numOfCol; i++) {
5,011✔
280
    if (!pFillCol[i].notFillCol) {
4,718✔
281
      int32_t          slotId = GET_DEST_SLOT_ID(pFillCol + i);
4,425✔
282
      SResultCellData* pECell = getResultCell(pEndRow, slotId);
4,425✔
283
      SPoint*          pPoint = taosArrayGet(pEndPoins, slotId);
4,424✔
284
      pPoint->key = pEndRow->key;
4,424✔
285
      memcpy(pPoint->val, pECell->pData, pECell->bytes);
4,424✔
286
    }
287
  }
288
}
293✔
289

290
static void setFillInfoStart(TSKEY ts, SInterval* pInterval, SStreamFillInfo* pFillInfo) {
1,399✔
291
  ts = taosTimeAdd(ts, pInterval->sliding, pInterval->slidingUnit, pInterval->precision, NULL);
1,399✔
292
  pFillInfo->start = ts;
1,399✔
293
}
1,399✔
294

295
static void setFillInfoEnd(TSKEY ts, SInterval* pInterval, SStreamFillInfo* pFillInfo) {
1,399✔
296
  ts = taosTimeAdd(ts, pInterval->sliding * -1, pInterval->slidingUnit, pInterval->precision, NULL);
1,399✔
297
  pFillInfo->end = ts;
1,399✔
298
}
1,399✔
299

300
static void setFillKeyInfo(TSKEY start, TSKEY end, SInterval* pInterval, SStreamFillInfo* pFillInfo) {
1,399✔
301
  setFillInfoStart(start, pInterval, pFillInfo);
1,399✔
302
  pFillInfo->current = pFillInfo->start;
1,399✔
303
  setFillInfoEnd(end, pInterval, pFillInfo);
1,399✔
304
}
1,399✔
305

306
void setDeleteFillValueInfo(TSKEY start, TSKEY end, SStreamFillSupporter* pFillSup, SStreamFillInfo* pFillInfo) {
200✔
307
  if (!hasPrevWindow(pFillSup) || !hasNextWindow(pFillSup)) {
200✔
308
    pFillInfo->needFill = false;
90✔
309
    return;
90✔
310
  }
311

312
  TSKEY realStart = taosTimeAdd(pFillSup->prev.key, pFillSup->interval.sliding, pFillSup->interval.slidingUnit,
110✔
313
                                pFillSup->interval.precision, NULL);
110✔
314

315
  pFillInfo->needFill = true;
110✔
316
  pFillInfo->start = realStart;
110✔
317
  pFillInfo->current = pFillInfo->start;
110✔
318
  pFillInfo->end = end;
110✔
319
  pFillInfo->pos = FILL_POS_INVALID;
110✔
320
  switch (pFillInfo->type) {
110!
321
    case TSDB_FILL_NULL:
44✔
322
    case TSDB_FILL_NULL_F:
323
    case TSDB_FILL_SET_VALUE:
324
    case TSDB_FILL_SET_VALUE_F:
325
      break;
44✔
326
    case TSDB_FILL_PREV:
22✔
327
      pFillInfo->pResRow = &pFillSup->prev;
22✔
328
      break;
22✔
329
    case TSDB_FILL_NEXT:
22✔
330
      pFillInfo->pResRow = &pFillSup->next;
22✔
331
      break;
22✔
332
    case TSDB_FILL_LINEAR: {
22✔
333
      setFillKeyInfo(pFillSup->prev.key, pFillSup->next.key, &pFillSup->interval, pFillInfo);
22✔
334
      pFillInfo->pLinearInfo->hasNext = false;
22✔
335
      pFillInfo->pLinearInfo->nextEnd = INT64_MIN;
22✔
336
      calcRowDeltaData(&pFillSup->next, pFillInfo->pLinearInfo->pEndPoints, pFillSup->pAllColInfo,
22✔
337
                       pFillSup->numOfAllCols);
338
      pFillInfo->pResRow = &pFillSup->prev;
22✔
339
      pFillInfo->pLinearInfo->winIndex = 0;
22✔
340
    } break;
22✔
341
    default:
×
342
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR));
×
343
      break;
×
344
  }
345
}
346

347
void copyNotFillExpData(SStreamFillSupporter* pFillSup, SStreamFillInfo* pFillInfo) {
745✔
348
  for (int32_t i = pFillSup->numOfFillCols; i < pFillSup->numOfAllCols; ++i) {
1,613✔
349
    SFillColInfo*    pFillCol = pFillSup->pAllColInfo + i;
868✔
350
    int32_t          slotId = GET_DEST_SLOT_ID(pFillCol);
868✔
351
    SResultCellData* pCell = getResultCell(pFillInfo->pResRow, slotId);
868✔
352
    SResultCellData* pCurCell = getResultCell(&pFillSup->cur, slotId);
868✔
353
    pCell->isNull = pCurCell->isNull;
868✔
354
    if (!pCurCell->isNull) {
868!
355
      memcpy(pCell->pData, pCurCell->pData, pCell->bytes);
868✔
356
    }
357
  }
358
}
745✔
359

360
void setFillValueInfo(SSDataBlock* pBlock, TSKEY ts, int32_t rowId, SStreamFillSupporter* pFillSup,
1,616✔
361
                      SStreamFillInfo* pFillInfo) {
362
  pFillInfo->preRowKey = pFillSup->cur.key;
1,616✔
363
  if (!hasPrevWindow(pFillSup) && !hasNextWindow(pFillSup)) {
1,616✔
364
    pFillInfo->needFill = false;
276✔
365
    pFillInfo->pos = FILL_POS_START;
276✔
366
    return;
276✔
367
  }
368
  TSKEY prevWKey = INT64_MIN;
1,340✔
369
  TSKEY nextWKey = INT64_MIN;
1,340✔
370
  if (hasPrevWindow(pFillSup)) {
1,340✔
371
    prevWKey = pFillSup->prev.key;
701✔
372
  }
373
  if (hasNextWindow(pFillSup)) {
1,340✔
374
    nextWKey = pFillSup->next.key;
902✔
375
  }
376

377
  pFillInfo->needFill = true;
1,340✔
378
  pFillInfo->pos = FILL_POS_INVALID;
1,340✔
379
  switch (pFillInfo->type) {
1,340!
380
    case TSDB_FILL_NULL:
574✔
381
    case TSDB_FILL_NULL_F:
382
    case TSDB_FILL_SET_VALUE:
383
    case TSDB_FILL_SET_VALUE_F: {
384
      if (pFillSup->prev.key == pFillInfo->preRowKey) {
574!
385
        resetFillWindow(&pFillSup->prev);
×
386
      }
387
      if (hasPrevWindow(pFillSup) && hasNextWindow(pFillSup)) {
574✔
388
        if (pFillSup->next.key == pFillInfo->nextRowKey) {
130✔
389
          pFillInfo->preRowKey = INT64_MIN;
128✔
390
          setFillKeyInfo(prevWKey, ts, &pFillSup->interval, pFillInfo);
128✔
391
          pFillInfo->pos = FILL_POS_END;
128✔
392
        } else {
393
          pFillInfo->needFill = false;
2✔
394
          pFillInfo->pos = FILL_POS_START;
2✔
395
        }
396
      } else if (hasPrevWindow(pFillSup)) {
444✔
397
        setFillKeyInfo(prevWKey, ts, &pFillSup->interval, pFillInfo);
188✔
398
        pFillInfo->pos = FILL_POS_END;
188✔
399
      } else {
400
        setFillKeyInfo(ts, nextWKey, &pFillSup->interval, pFillInfo);
256✔
401
        pFillInfo->pos = FILL_POS_START;
256✔
402
      }
403
      copyNotFillExpData(pFillSup, pFillInfo);
574✔
404
    } break;
574✔
405
    case TSDB_FILL_PREV: {
275✔
406
      if (hasNextWindow(pFillSup) && ((pFillSup->next.key != pFillInfo->nextRowKey) ||
275✔
407
                                      (pFillSup->next.key == pFillInfo->nextRowKey && hasNextNextWindow(pFillSup)) ||
149!
408
                                      (pFillSup->next.key == pFillInfo->nextRowKey && !hasPrevWindow(pFillSup)))) {
122!
409
        setFillKeyInfo(ts, nextWKey, &pFillSup->interval, pFillInfo);
136✔
410
        pFillInfo->pos = FILL_POS_START;
136✔
411
        resetFillWindow(&pFillSup->prev);
136✔
412
        pFillSup->prev.key = pFillSup->cur.key;
136✔
413
        pFillSup->prev.pRowVal = pFillSup->cur.pRowVal;
136✔
414
      } else if (hasPrevWindow(pFillSup)) {
139!
415
        setFillKeyInfo(prevWKey, ts, &pFillSup->interval, pFillInfo);
139✔
416
        pFillInfo->pos = FILL_POS_END;
139✔
417
        pFillInfo->preRowKey = INT64_MIN;
139✔
418
      }
419
      pFillInfo->pResRow = &pFillSup->prev;
275✔
420
    } break;
275✔
421
    case TSDB_FILL_NEXT: {
258✔
422
      if (hasPrevWindow(pFillSup)) {
258✔
423
        setFillKeyInfo(prevWKey, ts, &pFillSup->interval, pFillInfo);
147✔
424
        pFillInfo->pos = FILL_POS_END;
147✔
425
        resetFillWindow(&pFillSup->next);
147✔
426
        pFillSup->next.key = pFillSup->cur.key;
147✔
427
        pFillSup->next.pRowVal = pFillSup->cur.pRowVal;
147✔
428
        pFillInfo->preRowKey = INT64_MIN;
147✔
429
      } else {
430
        setFillKeyInfo(ts, nextWKey, &pFillSup->interval, pFillInfo);
111✔
431
        pFillInfo->pos = FILL_POS_START;
111✔
432
      }
433
      pFillInfo->pResRow = &pFillSup->next;
258✔
434
    } break;
258✔
435
    case TSDB_FILL_LINEAR: {
233✔
436
      pFillInfo->pLinearInfo->winIndex = 0;
233✔
437
      if (hasPrevWindow(pFillSup) && hasNextWindow(pFillSup)) {
233✔
438
        setFillKeyInfo(prevWKey, ts, &pFillSup->interval, pFillInfo);
39✔
439
        pFillInfo->pos = FILL_POS_MID;
39✔
440
        pFillInfo->pLinearInfo->nextEnd = nextWKey;
39✔
441
        calcRowDeltaData(&pFillSup->cur, pFillInfo->pLinearInfo->pEndPoints, pFillSup->pAllColInfo,
39✔
442
                         pFillSup->numOfAllCols);
443
        pFillInfo->pResRow = &pFillSup->prev;
39✔
444

445
        calcRowDeltaData(&pFillSup->next, pFillInfo->pLinearInfo->pNextEndPoints, pFillSup->pAllColInfo,
39✔
446
                         pFillSup->numOfAllCols);
447
        pFillInfo->pLinearInfo->hasNext = true;
39✔
448
      } else if (hasPrevWindow(pFillSup)) {
194✔
449
        setFillKeyInfo(prevWKey, ts, &pFillSup->interval, pFillInfo);
55✔
450
        pFillInfo->pos = FILL_POS_END;
55✔
451
        pFillInfo->pLinearInfo->nextEnd = INT64_MIN;
55✔
452
        calcRowDeltaData(&pFillSup->cur, pFillInfo->pLinearInfo->pEndPoints, pFillSup->pAllColInfo,
55✔
453
                         pFillSup->numOfAllCols);
454
        pFillInfo->pResRow = &pFillSup->prev;
55✔
455
        pFillInfo->pLinearInfo->hasNext = false;
55✔
456
      } else {
457
        setFillKeyInfo(ts, nextWKey, &pFillSup->interval, pFillInfo);
139✔
458
        pFillInfo->pos = FILL_POS_START;
139✔
459
        pFillInfo->pLinearInfo->nextEnd = INT64_MIN;
139✔
460
        calcRowDeltaData(&pFillSup->next, pFillInfo->pLinearInfo->pEndPoints, pFillSup->pAllColInfo,
139✔
461
                         pFillSup->numOfAllCols);
462
        pFillInfo->pResRow = &pFillSup->cur;
139✔
463
        pFillInfo->pLinearInfo->hasNext = false;
139✔
464
      }
465
    } break;
233✔
466
    default:
×
467
      qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR));
×
468
      break;
×
469
  }
470
}
471

472
int32_t checkResult(SStreamFillSupporter* pFillSup, TSKEY ts, uint64_t groupId, bool* pRes) {
203,083✔
473
  int32_t code = TSDB_CODE_SUCCESS;
203,083✔
474
  int32_t lino = 0;
203,083✔
475
  SWinKey key = {.groupId = groupId, .ts = ts};
203,083✔
476
  if (tSimpleHashGet(pFillSup->pResMap, &key, sizeof(SWinKey)) != NULL) {
203,083✔
477
    (*pRes) = false;
50,208✔
478
    goto _end;
50,208✔
479
  }
480
  code = tSimpleHashPut(pFillSup->pResMap, &key, sizeof(SWinKey), NULL, 0);
152,987✔
481
  QUERY_CHECK_CODE(code, lino, _end);
153,024!
482
  (*pRes) = true;
153,024✔
483

484
_end:
203,232✔
485
  if (code != TSDB_CODE_SUCCESS) {
203,232!
486
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
487
  }
488
  return code;
203,236✔
489
}
490

491
static int32_t buildFillResult(SResultRowData* pResRow, SStreamFillSupporter* pFillSup, TSKEY ts, SSDataBlock* pBlock,
37,327✔
492
                               bool* pRes, bool isFilld) {
493
  int32_t code = TSDB_CODE_SUCCESS;
37,327✔
494
  int32_t lino = 0;
37,327✔
495
  if (pBlock->info.rows >= pBlock->info.capacity) {
37,327✔
496
    (*pRes) = false;
4✔
497
    goto _end;
4✔
498
  }
499
  uint64_t groupId = pBlock->info.id.groupId;
37,323✔
500
  bool     ckRes = true;
37,323✔
501
  code = checkResult(pFillSup, ts, groupId, &ckRes);
37,323✔
502
  QUERY_CHECK_CODE(code, lino, _end);
37,325!
503

504
  if (pFillSup->hasDelete && !ckRes) {
37,325✔
505
    (*pRes) = true;
28✔
506
    goto _end;
28✔
507
  }
508
  for (int32_t i = 0; i < pFillSup->numOfAllCols; ++i) {
162,234✔
509
    SFillColInfo*    pFillCol = pFillSup->pAllColInfo + i;
124,938✔
510
    int32_t          slotId = GET_DEST_SLOT_ID(pFillCol);
124,938✔
511
    SColumnInfoData* pColData = taosArrayGet(pBlock->pDataBlock, slotId);
124,938✔
512
    SFillInfo        tmpInfo = {
124,930✔
513
               .currentKey = ts,
514
               .order = TSDB_ORDER_ASC,
515
               .interval = pFillSup->interval,
516
               .isFilled = isFilld,
517
    };
518
    bool filled = fillIfWindowPseudoColumn(&tmpInfo, pFillCol, pColData, pBlock->info.rows);
124,930✔
519
    if (!filled) {
125,402✔
520
      SResultCellData* pCell = getResultCell(pResRow, slotId);
88,225✔
521
      code = setRowCell(pColData, pBlock->info.rows, pCell);
87,766✔
522
      QUERY_CHECK_CODE(code, lino, _end);
87,760!
523
    }
524
  }
525
  pBlock->info.rows++;
37,296✔
526
  (*pRes) = true;
37,296✔
527

528
_end:
37,328✔
529
  if (code != TSDB_CODE_SUCCESS) {
37,328!
530
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
531
  }
532
  return code;
37,327✔
533
}
534

535
bool hasRemainCalc(SStreamFillInfo* pFillInfo) {
207,137✔
536
  if (pFillInfo->current != INT64_MIN && pFillInfo->current <= pFillInfo->end) {
207,137✔
537
    return true;
200,334✔
538
  }
539
  return false;
6,803✔
540
}
541

542
static void doStreamFillNormal(SStreamFillSupporter* pFillSup, SStreamFillInfo* pFillInfo, SSDataBlock* pBlock) {
1,156✔
543
  int32_t code = TSDB_CODE_SUCCESS;
1,156✔
544
  int32_t lino = 0;
1,156✔
545
  while (hasRemainCalc(pFillInfo) && pBlock->info.rows < pBlock->info.capacity) {
36,210✔
546
    STimeWindow st = {.skey = pFillInfo->current, .ekey = pFillInfo->current};
35,052✔
547
    if (inWinRange(&pFillSup->winRange, &st)) {
35,052!
548
      bool res = true;
35,051✔
549
      code = buildFillResult(pFillInfo->pResRow, pFillSup, pFillInfo->current, pBlock, &res, true);
35,051✔
550
      QUERY_CHECK_CODE(code, lino, _end);
35,054!
551
    }
552
    pFillInfo->current = taosTimeAdd(pFillInfo->current, pFillSup->interval.sliding, pFillSup->interval.slidingUnit,
35,054✔
553
                                     pFillSup->interval.precision, NULL);
35,054✔
554
  }
555

556
_end:
1,156✔
557
  if (code != TSDB_CODE_SUCCESS) {
1,156!
558
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
559
  }
560
}
1,156✔
561

562
static void doStreamFillLinear(SStreamFillSupporter* pFillSup, SStreamFillInfo* pFillInfo, SSDataBlock* pBlock) {
286✔
563
  int32_t code = TSDB_CODE_SUCCESS;
286✔
564
  int32_t lino = 0;
286✔
565
  while (hasRemainCalc(pFillInfo) && pBlock->info.rows < pBlock->info.capacity) {
14,842✔
566
    uint64_t    groupId = pBlock->info.id.groupId;
14,556✔
567
    SWinKey     key = {.groupId = groupId, .ts = pFillInfo->current};
14,556✔
568
    STimeWindow st = {.skey = pFillInfo->current, .ekey = pFillInfo->current};
14,556✔
569
    bool        ckRes = true;
14,556✔
570
    code = checkResult(pFillSup, pFillInfo->current, groupId, &ckRes);
14,556✔
571
    QUERY_CHECK_CODE(code, lino, _end);
14,556!
572

573
    if ((pFillSup->hasDelete && !ckRes) || !inWinRange(&pFillSup->winRange, &st)) {
14,556!
574
      pFillInfo->current = taosTimeAdd(pFillInfo->current, pFillSup->interval.sliding, pFillSup->interval.slidingUnit,
16✔
575
                                       pFillSup->interval.precision, NULL);
8✔
576
      pFillInfo->pLinearInfo->winIndex++;
8✔
577
      continue;
8✔
578
    }
579
    pFillInfo->pLinearInfo->winIndex++;
14,548✔
580
    for (int32_t i = 0; i < pFillSup->numOfAllCols; ++i) {
49,937✔
581
      SFillColInfo* pFillCol = pFillSup->pAllColInfo + i;
35,387✔
582
      SFillInfo     tmp = {
35,387✔
583
              .currentKey = pFillInfo->current,
35,387✔
584
              .order = TSDB_ORDER_ASC,
585
              .interval = pFillSup->interval,
586
              .isFilled = true,
587
      };
588

589
      int32_t          slotId = GET_DEST_SLOT_ID(pFillCol);
35,387✔
590
      SColumnInfoData* pColData = taosArrayGet(pBlock->pDataBlock, slotId);
35,387✔
591
      int16_t          type = pColData->info.type;
35,409✔
592
      SResultCellData* pCell = getResultCell(pFillInfo->pResRow, slotId);
35,409✔
593
      int32_t          index = pBlock->info.rows;
35,315✔
594
      if (pFillCol->notFillCol) {
35,315✔
595
        bool filled = fillIfWindowPseudoColumn(&tmp, pFillCol, pColData, index);
14,548✔
596
        if (!filled) {
14,548✔
597
          code = setRowCell(pColData, index, pCell);
5✔
598
          QUERY_CHECK_CODE(code, lino, _end);
×
599
        }
600
      } else {
601
        if (IS_VAR_DATA_TYPE(type) || type == TSDB_DATA_TYPE_BOOL || pCell->isNull) {
20,767!
602
          colDataSetNULL(pColData, index);
23!
603
          continue;
23✔
604
        }
605
        SPoint* pEnd = taosArrayGet(pFillInfo->pLinearInfo->pEndPoints, slotId);
20,744✔
606
        double  vCell = 0;
20,742✔
607
        SPoint  start = {0};
20,742✔
608
        start.key = pFillInfo->pResRow->key;
20,742✔
609
        start.val = pCell->pData;
20,742✔
610

611
        SPoint cur = {0};
20,742✔
612
        cur.key = pFillInfo->current;
20,742✔
613
        cur.val = taosMemoryCalloc(1, pCell->bytes);
20,742!
614
        QUERY_CHECK_NULL(cur.val, code, lino, _end, terrno);
20,817!
615
        taosGetLinearInterpolationVal(&cur, pCell->type, &start, pEnd, pCell->type, typeGetTypeModFromColInfo(&pColData->info));
20,817✔
616
        code = colDataSetVal(pColData, index, (const char*)cur.val, false);
20,781✔
617
        QUERY_CHECK_CODE(code, lino, _end);
20,746!
618
        destroySPoint(&cur);
20,746✔
619
      }
620
    }
621
    pFillInfo->current = taosTimeAdd(pFillInfo->current, pFillSup->interval.sliding, pFillSup->interval.slidingUnit,
29,098✔
622
                                     pFillSup->interval.precision, NULL);
14,550✔
623
    pBlock->info.rows++;
14,548✔
624
  }
625

626
_end:
286✔
627
  if (code != TSDB_CODE_SUCCESS) {
286!
628
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
629
  }
630
}
286✔
631

632
static void keepResultInDiscBuf(SOperatorInfo* pOperator, uint64_t groupId, SResultRowData* pRow, int32_t len) {
1,616✔
633
  SStorageAPI* pAPI = &pOperator->pTaskInfo->storageAPI;
1,616✔
634

635
  SWinKey key = {.groupId = groupId, .ts = pRow->key};
1,616✔
636
  int32_t code = pAPI->stateStore.streamStateFillPut(pOperator->pTaskInfo->streamInfo.pState, &key, pRow->pRowVal, len);
1,616✔
637
  qDebug("===stream===fill operator save key ts:%" PRId64 " group id:%" PRIu64 "  code:%d", key.ts, key.groupId, code);
1,616✔
638
  if (code != TSDB_CODE_SUCCESS) {
1,616!
639
    qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(code));
×
640
  }
641
}
1,616✔
642

643
void doStreamFillRange(SStreamFillInfo* pFillInfo, SStreamFillSupporter* pFillSup, SSDataBlock* pRes) {
1,681✔
644
  int32_t code = TSDB_CODE_SUCCESS;
1,681✔
645
  int32_t lino = 0;
1,681✔
646
  bool    res = false;
1,681✔
647
  if (pFillInfo->needFill == false) {
1,681✔
648
    code = buildFillResult(&pFillSup->cur, pFillSup, pFillSup->cur.key, pRes, &res, false);
278✔
649
    QUERY_CHECK_CODE(code, lino, _end);
278!
650
    return;
278✔
651
  }
652

653
  if (pFillInfo->pos == FILL_POS_START) {
1,403✔
654
    code = buildFillResult(&pFillSup->cur, pFillSup, pFillSup->cur.key, pRes, &res, false);
642✔
655
    QUERY_CHECK_CODE(code, lino, _end);
642!
656
    if (res) {
642!
657
      pFillInfo->pos = FILL_POS_INVALID;
642✔
658
    }
659
  }
660
  if (pFillInfo->type != TSDB_FILL_LINEAR) {
1,403✔
661
    doStreamFillNormal(pFillSup, pFillInfo, pRes);
1,156✔
662
  } else {
663
    doStreamFillLinear(pFillSup, pFillInfo, pRes);
247✔
664

665
    if (pFillInfo->pos == FILL_POS_MID) {
247✔
666
      code = buildFillResult(&pFillSup->cur, pFillSup, pFillSup->cur.key, pRes, &res, false);
39✔
667
      QUERY_CHECK_CODE(code, lino, _end);
39!
668
      if (res) {
39!
669
        pFillInfo->pos = FILL_POS_INVALID;
39✔
670
      }
671
    }
672

673
    if (pFillInfo->current > pFillInfo->end && pFillInfo->pLinearInfo->hasNext) {
247✔
674
      pFillInfo->pLinearInfo->hasNext = false;
39✔
675
      pFillInfo->pLinearInfo->winIndex = 0;
39✔
676
      taosArraySwap(pFillInfo->pLinearInfo->pEndPoints, pFillInfo->pLinearInfo->pNextEndPoints);
39✔
677
      pFillInfo->pResRow = &pFillSup->cur;
39✔
678
      setFillKeyInfo(pFillSup->cur.key, pFillInfo->pLinearInfo->nextEnd, &pFillSup->interval, pFillInfo);
39✔
679
      doStreamFillLinear(pFillSup, pFillInfo, pRes);
39✔
680
    }
681
  }
682
  if (pFillInfo->pos == FILL_POS_END) {
1,403✔
683
    code = buildFillResult(&pFillSup->cur, pFillSup, pFillSup->cur.key, pRes, &res, false);
661✔
684
    QUERY_CHECK_CODE(code, lino, _end);
661!
685
    if (res) {
661✔
686
      pFillInfo->pos = FILL_POS_INVALID;
657✔
687
    }
688
  }
689

690
_end:
746✔
691
  if (code != TSDB_CODE_SUCCESS) {
1,403!
692
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
693
  }
694
}
695

696
int32_t keepBlockRowInDiscBuf(SOperatorInfo* pOperator, SStreamFillInfo* pFillInfo, SSDataBlock* pBlock, TSKEY* tsCol,
1,616✔
697
                           int32_t rowId, uint64_t groupId, int32_t rowSize) {
698
  int32_t code = TSDB_CODE_SUCCESS;
1,616✔
699
  int32_t lino = 0;
1,616✔
700
  TSKEY ts = tsCol[rowId];
1,616✔
701
  pFillInfo->nextRowKey = ts;
1,616✔
702
  SResultRowData tmpNextRow = {.key = ts};
1,616✔
703
  tmpNextRow.pRowVal = taosMemoryCalloc(1, rowSize);
1,616!
704
  QUERY_CHECK_NULL(tmpNextRow.pRowVal, code, lino, _end, terrno);
1,616!
705
  transBlockToResultRow(pBlock, rowId, ts, &tmpNextRow);
1,616✔
706
  keepResultInDiscBuf(pOperator, groupId, &tmpNextRow, rowSize);
1,616✔
707
  taosMemoryFreeClear(tmpNextRow.pRowVal);
1,616!
708

709
_end:
×
710
  if (code != TSDB_CODE_SUCCESS) {
1,616!
711
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
712
  }
713
  return code;
1,616✔
714
}
715

716
static void doFillResults(SOperatorInfo* pOperator, SStreamFillSupporter* pFillSup, SStreamFillInfo* pFillInfo,
1,615✔
717
                          SSDataBlock* pBlock, TSKEY* tsCol, int32_t rowId, SSDataBlock* pRes) {
718
  uint64_t groupId = pBlock->info.id.groupId;
1,615✔
719
  getWindowFromDiscBuf(pOperator, tsCol[rowId], groupId, pFillSup);
1,615✔
720
  if (pFillSup->prev.key == pFillInfo->preRowKey) {
1,616✔
721
    resetFillWindow(&pFillSup->prev);
915✔
722
  }
723
  setFillValueInfo(pBlock, tsCol[rowId], rowId, pFillSup, pFillInfo);
1,616✔
724
  doStreamFillRange(pFillInfo, pFillSup, pRes);
1,616✔
725
}
1,616✔
726

727
static void doStreamFillImpl(SOperatorInfo* pOperator) {
730✔
728
  int32_t                  code = TSDB_CODE_SUCCESS;
730✔
729
  int32_t                  lino = 0;
730✔
730
  SStreamFillOperatorInfo* pInfo = pOperator->info;
730✔
731
  SExecTaskInfo*           pTaskInfo = pOperator->pTaskInfo;
730✔
732
  SStreamFillSupporter*    pFillSup = pInfo->pFillSup;
730✔
733
  SStreamFillInfo*         pFillInfo = pInfo->pFillInfo;
730✔
734
  SSDataBlock*             pBlock = pInfo->pSrcBlock;
730✔
735
  uint64_t                 groupId = pBlock->info.id.groupId;
730✔
736
  SSDataBlock*             pRes = pInfo->pRes;
730✔
737
  SColumnInfoData*         pTsCol = taosArrayGet(pInfo->pSrcBlock->pDataBlock, pInfo->primaryTsCol);
730✔
738
  TSKEY*                   tsCol = (TSKEY*)pTsCol->pData;
730✔
739
  pRes->info.id.groupId = groupId;
730✔
740
  pInfo->srcRowIndex++;
730✔
741

742
  if (pInfo->srcRowIndex == 0) {
730✔
743
    code = keepBlockRowInDiscBuf(pOperator, pFillInfo, pBlock, tsCol, pInfo->srcRowIndex, groupId, pFillSup->rowSize);
727✔
744
    QUERY_CHECK_CODE(code, lino, _end);
727!
745
    pInfo->srcRowIndex++;
727✔
746
  }
747

748
  while (pInfo->srcRowIndex < pBlock->info.rows) {
1,616✔
749
    code = keepBlockRowInDiscBuf(pOperator, pFillInfo, pBlock, tsCol, pInfo->srcRowIndex, groupId, pFillSup->rowSize);
889✔
750
    QUERY_CHECK_CODE(code, lino, _end);
889!
751
    doFillResults(pOperator, pFillSup, pFillInfo, pBlock, tsCol, pInfo->srcRowIndex - 1, pRes);
889✔
752
    if (pInfo->pRes->info.rows == pInfo->pRes->info.capacity) {
889✔
753
      code = blockDataUpdateTsWindow(pRes, pInfo->primaryTsCol);
3✔
754
      QUERY_CHECK_CODE(code, lino, _end);
3!
755
      return;
3✔
756
    }
757
    pInfo->srcRowIndex++;
886✔
758
  }
759
  doFillResults(pOperator, pFillSup, pFillInfo, pBlock, tsCol, pInfo->srcRowIndex - 1, pRes);
727✔
760
  code = blockDataUpdateTsWindow(pRes, pInfo->primaryTsCol);
727✔
761
  QUERY_CHECK_CODE(code, lino, _end);
727!
762
  blockDataCleanup(pInfo->pSrcBlock);
727✔
763

764
_end:
727✔
765
  if (code != TSDB_CODE_SUCCESS) {
727!
766
    qError("%s failed at line %d since %s. task:%s", __func__, lino, tstrerror(code), GET_TASKID(pTaskInfo));
×
767
  }
768
}
769

770
static int32_t buildDeleteRange(SOperatorInfo* pOp, TSKEY start, TSKEY end, uint64_t groupId, SSDataBlock* delRes) {
90✔
771
  int32_t          code = TSDB_CODE_SUCCESS;
90✔
772
  int32_t          lino = 0;
90✔
773
  SStorageAPI*     pAPI = &pOp->pTaskInfo->storageAPI;
90✔
774
  void*            pState = pOp->pTaskInfo->streamInfo.pState;
90✔
775
  SExecTaskInfo*   pTaskInfo = pOp->pTaskInfo;
90✔
776
  SSDataBlock*     pBlock = delRes;
90✔
777
  SColumnInfoData* pStartCol = taosArrayGet(pBlock->pDataBlock, START_TS_COLUMN_INDEX);
90✔
778
  SColumnInfoData* pEndCol = taosArrayGet(pBlock->pDataBlock, END_TS_COLUMN_INDEX);
90✔
779
  SColumnInfoData* pUidCol = taosArrayGet(pBlock->pDataBlock, UID_COLUMN_INDEX);
90✔
780
  SColumnInfoData* pGroupCol = taosArrayGet(pBlock->pDataBlock, GROUPID_COLUMN_INDEX);
90✔
781
  SColumnInfoData* pCalStartCol = taosArrayGet(pBlock->pDataBlock, CALCULATE_START_TS_COLUMN_INDEX);
90✔
782
  SColumnInfoData* pCalEndCol = taosArrayGet(pBlock->pDataBlock, CALCULATE_END_TS_COLUMN_INDEX);
90✔
783
  SColumnInfoData* pTbNameCol = taosArrayGet(pBlock->pDataBlock, TABLE_NAME_COLUMN_INDEX);
90✔
784
  code = colDataSetVal(pStartCol, pBlock->info.rows, (const char*)&start, false);
90✔
785
  QUERY_CHECK_CODE(code, lino, _end);
90!
786

787
  code = colDataSetVal(pEndCol, pBlock->info.rows, (const char*)&end, false);
90✔
788
  QUERY_CHECK_CODE(code, lino, _end);
90!
789

790
  colDataSetNULL(pUidCol, pBlock->info.rows);
90!
791
  code = colDataSetVal(pGroupCol, pBlock->info.rows, (const char*)&groupId, false);
90✔
792
  QUERY_CHECK_CODE(code, lino, _end);
90!
793

794
  colDataSetNULL(pCalStartCol, pBlock->info.rows);
90!
795
  colDataSetNULL(pCalEndCol, pBlock->info.rows);
90!
796

797
  SColumnInfoData* pTableCol = taosArrayGet(pBlock->pDataBlock, TABLE_NAME_COLUMN_INDEX);
90✔
798

799
  void*   tbname = NULL;
90✔
800
  int32_t winCode = TSDB_CODE_SUCCESS;
90✔
801
  code = pAPI->stateStore.streamStateGetParName(pOp->pTaskInfo->streamInfo.pState, groupId, &tbname, false, &winCode);
90✔
802
  QUERY_CHECK_CODE(code, lino, _end);
90!
803
  if (winCode != TSDB_CODE_SUCCESS) {
90✔
804
    colDataSetNULL(pTableCol, pBlock->info.rows);
15!
805
  } else {
806
    char parTbName[VARSTR_HEADER_SIZE + TSDB_TABLE_NAME_LEN];
807
    STR_WITH_MAXSIZE_TO_VARSTR(parTbName, tbname, sizeof(parTbName));
75✔
808
    code = colDataSetVal(pTableCol, pBlock->info.rows, (const char*)parTbName, false);
75✔
809
    QUERY_CHECK_CODE(code, lino, _end);
75!
810
    pAPI->stateStore.streamStateFreeVal(tbname);
75✔
811
  }
812

813
  pBlock->info.rows++;
90✔
814

815
_end:
90✔
816
  if (code != TSDB_CODE_SUCCESS) {
90!
817
    qError("%s failed at line %d since %s. task:%s", __func__, lino, tstrerror(code), GET_TASKID(pTaskInfo));
×
818
  }
819
  return code;
90✔
820
}
821

822
int32_t buildDeleteResult(SOperatorInfo* pOperator, TSKEY startTs, TSKEY endTs, uint64_t groupId, SSDataBlock* delRes) {
90✔
823
  int32_t                  code = TSDB_CODE_SUCCESS;
90✔
824
  int32_t                  lino = 0;
90✔
825
  SStreamFillOperatorInfo* pInfo = pOperator->info;
90✔
826
  SStreamFillSupporter*    pFillSup = pInfo->pFillSup;
90✔
827
  SExecTaskInfo*           pTaskInfo = pOperator->pTaskInfo;
90✔
828
  if (hasPrevWindow(pFillSup)) {
90✔
829
    TSKEY start = getNextWindowTs(pFillSup->prev.key, &pFillSup->interval);
55✔
830
    code = buildDeleteRange(pOperator, start, endTs, groupId, delRes);
55✔
831
    QUERY_CHECK_CODE(code, lino, _end);
55!
832
  } else if (hasNextWindow(pFillSup)) {
35✔
833
    TSKEY end = getPrevWindowTs(pFillSup->next.key, &pFillSup->interval);
30✔
834
    code = buildDeleteRange(pOperator, startTs, end, groupId, delRes);
30✔
835
    QUERY_CHECK_CODE(code, lino, _end);
30!
836
  } else {
837
    code = buildDeleteRange(pOperator, startTs, endTs, groupId, delRes);
5✔
838
    QUERY_CHECK_CODE(code, lino, _end);
5!
839
  }
840

841
_end:
5✔
842
  if (code != TSDB_CODE_SUCCESS) {
90!
843
    qError("%s failed at line %d since %s. task:%s", __func__, lino, tstrerror(code), GET_TASKID(pTaskInfo));
×
844
  }
845
  return code;
90✔
846
}
847

848
static int32_t doDeleteFillResultImpl(SOperatorInfo* pOperator, TSKEY startTs, TSKEY endTs, uint64_t groupId) {
145✔
849
  int32_t                  code = TSDB_CODE_SUCCESS;
145✔
850
  int32_t                  lino = 0;
145✔
851
  SStorageAPI*             pAPI = &pOperator->pTaskInfo->storageAPI;
145✔
852
  SStreamFillOperatorInfo* pInfo = pOperator->info;
145✔
853
  SExecTaskInfo*           pTaskInfo = pOperator->pTaskInfo;
145✔
854
  getWindowFromDiscBuf(pOperator, startTs, groupId, pInfo->pFillSup);
145✔
855
  setDeleteFillValueInfo(startTs, endTs, pInfo->pFillSup, pInfo->pFillInfo);
145✔
856
  SWinKey key = {.ts = startTs, .groupId = groupId};
145✔
857
  pAPI->stateStore.streamStateFillDel(pOperator->pTaskInfo->streamInfo.pState, &key);
145✔
858
  if (!pInfo->pFillInfo->needFill) {
145✔
859
    code = buildDeleteResult(pOperator, startTs, endTs, groupId, pInfo->pDelRes);
90✔
860
    QUERY_CHECK_CODE(code, lino, _end);
90!
861
  } else {
862
    STimeFillRange tw = {
55✔
863
        .skey = startTs,
864
        .ekey = endTs,
865
        .groupId = groupId,
866
    };
867
    void* tmp = taosArrayPush(pInfo->pFillInfo->delRanges, &tw);
55✔
868
    if (!tmp) {
55!
869
      code = terrno;
×
870
      QUERY_CHECK_CODE(code, lino, _end);
×
871
    }
872
  }
873

874
_end:
145✔
875
  if (code != TSDB_CODE_SUCCESS) {
145!
876
    qError("%s failed at line %d since %s. task:%s", __func__, lino, tstrerror(code), GET_TASKID(pTaskInfo));
×
877
  }
878
  return code;
145✔
879
}
880

881
static void getWindowInfoByKey(SStorageAPI* pAPI, void* pState, TSKEY ts, int64_t groupId, SResultRowData* pWinData) {
×
882
  SWinKey key = {.ts = ts, .groupId = groupId};
×
883
  void*   val = NULL;
×
884
  int32_t len = 0;
×
885
  int32_t code = pAPI->stateStore.streamStateFillGet(pState, &key, (void**)&val, &len, NULL);
×
886
  if (code != TSDB_CODE_SUCCESS) {
×
887
    qDebug("get window info by key failed, Data may be deleted, try next window. ts:%" PRId64 ", groupId:%" PRId64, ts,
×
888
           groupId);
889
    SStreamStateCur* pCur = pAPI->stateStore.streamStateFillSeekKeyNext(pState, &key);
×
890
    code = pAPI->stateStore.streamStateFillGetGroupKVByCur(pCur, &key, (const void**)&val, &len);
×
891
    pAPI->stateStore.streamStateFreeCur(pCur);
×
892
    qDebug("get window info by key ts:%" PRId64 ", groupId:%" PRId64 ", res%d", ts, groupId, code);
×
893
  }
894

895
  if (code == TSDB_CODE_SUCCESS) {
×
896
    resetFillWindow(pWinData);
×
897
    pWinData->key = key.ts;
×
898
    pWinData->pRowVal = val;
×
899
  }
900
}
×
901

902
static void doDeleteFillFinalize(SOperatorInfo* pOperator) {
2,193✔
903
  SStorageAPI* pAPI = &pOperator->pTaskInfo->storageAPI;
2,193✔
904

905
  SStreamFillOperatorInfo* pInfo = pOperator->info;
2,193✔
906
  SStreamFillInfo*         pFillInfo = pInfo->pFillInfo;
2,193✔
907
  int32_t                  size = taosArrayGetSize(pFillInfo->delRanges);
2,193✔
908
  while (pFillInfo->delIndex < size) {
2,250✔
909
    STimeFillRange* range = taosArrayGet(pFillInfo->delRanges, pFillInfo->delIndex);
55✔
910
    if (pInfo->pRes->info.id.groupId != 0 && pInfo->pRes->info.id.groupId != range->groupId) {
55!
911
      return;
×
912
    }
913
    getWindowFromDiscBuf(pOperator, range->skey, range->groupId, pInfo->pFillSup);
55✔
914
    TSKEY realEnd = range->ekey + 1;
55✔
915
    if (pInfo->pFillInfo->type == TSDB_FILL_NEXT && pInfo->pFillSup->next.key != realEnd) {
55!
916
      getWindowInfoByKey(pAPI, pOperator->pTaskInfo->streamInfo.pState, realEnd, range->groupId,
×
917
                         &pInfo->pFillSup->next);
×
918
    }
919
    setDeleteFillValueInfo(range->skey, range->ekey, pInfo->pFillSup, pInfo->pFillInfo);
55✔
920
    pFillInfo->delIndex++;
55✔
921
    if (pInfo->pFillInfo->needFill) {
55!
922
      doStreamFillRange(pInfo->pFillInfo, pInfo->pFillSup, pInfo->pRes);
55✔
923
      pInfo->pRes->info.id.groupId = range->groupId;
55✔
924
    }
925
  }
926
}
927

928
static int32_t doDeleteFillResult(SOperatorInfo* pOperator) {
135✔
929
  int32_t                  code = TSDB_CODE_SUCCESS;
135✔
930
  int32_t                  lino = 0;
135✔
931
  SStorageAPI*             pAPI = &pOperator->pTaskInfo->storageAPI;
135✔
932
  SStreamFillOperatorInfo* pInfo = pOperator->info;
135✔
933
  SStreamFillInfo*         pFillInfo = pInfo->pFillInfo;
135✔
934
  SSDataBlock*             pBlock = pInfo->pSrcDelBlock;
135✔
935
  SExecTaskInfo*           pTaskInfo = pOperator->pTaskInfo;
135✔
936

937
  SColumnInfoData* pStartCol = taosArrayGet(pBlock->pDataBlock, START_TS_COLUMN_INDEX);
135✔
938
  TSKEY*           tsStarts = (TSKEY*)pStartCol->pData;
135✔
939
  SColumnInfoData* pGroupCol = taosArrayGet(pBlock->pDataBlock, GROUPID_COLUMN_INDEX);
135✔
940
  uint64_t*        groupIds = (uint64_t*)pGroupCol->pData;
135✔
941
  while (pInfo->srcDelRowIndex < pBlock->info.rows) {
340✔
942
    TSKEY            ts = tsStarts[pInfo->srcDelRowIndex];
205✔
943
    TSKEY            endTs = ts;
205✔
944
    uint64_t         groupId = groupIds[pInfo->srcDelRowIndex];
205✔
945
    SWinKey          key = {.ts = ts, .groupId = groupId};
205✔
946
    SStreamStateCur* pCur = pAPI->stateStore.streamStateGetAndCheckCur(pOperator->pTaskInfo->streamInfo.pState, &key);
205✔
947

948
    if (!pCur) {
205✔
949
      pInfo->srcDelRowIndex++;
60✔
950
      continue;
60✔
951
    }
952

953
    SWinKey nextKey = {.groupId = groupId, .ts = ts};
145✔
954
    while (pInfo->srcDelRowIndex < pBlock->info.rows) {
406✔
955
      TSKEY    delTs = tsStarts[pInfo->srcDelRowIndex];
331✔
956
      uint64_t delGroupId = groupIds[pInfo->srcDelRowIndex];
331✔
957
      int32_t  winCode = TSDB_CODE_SUCCESS;
331✔
958
      if (groupId != delGroupId) {
331!
959
        break;
70✔
960
      }
961
      if (delTs > nextKey.ts) {
331✔
962
        break;
10✔
963
      }
964

965
      SWinKey delKey = {.groupId = delGroupId, .ts = delTs};
321✔
966
      if (delTs == nextKey.ts) {
321✔
967
        pAPI->stateStore.streamStateCurNext(pOperator->pTaskInfo->streamInfo.pState, pCur);
243✔
968
        winCode = pAPI->stateStore.streamStateFillGetGroupKVByCur(pCur, &nextKey, NULL, NULL);
243✔
969
        // ts will be deleted later
970
        if (delTs != ts) {
243✔
971
          pAPI->stateStore.streamStateFillDel(pOperator->pTaskInfo->streamInfo.pState, &delKey);
98✔
972
          pAPI->stateStore.streamStateFreeCur(pCur);
98✔
973
          pCur = pAPI->stateStore.streamStateGetAndCheckCur(pOperator->pTaskInfo->streamInfo.pState, &nextKey);
98✔
974
        }
975
        endTs = TMAX(delTs, nextKey.ts - 1);
243✔
976
        if (winCode != TSDB_CODE_SUCCESS) {
243✔
977
          break;
60✔
978
        }
979
      }
980
      pInfo->srcDelRowIndex++;
261✔
981
    }
982

983
    pAPI->stateStore.streamStateFreeCur(pCur);
145✔
984
    code = doDeleteFillResultImpl(pOperator, ts, endTs, groupId);
145✔
985
    QUERY_CHECK_CODE(code, lino, _end);
145!
986
  }
987

988
  pFillInfo->current = pFillInfo->end + 1;
135✔
989

990
_end:
135✔
991
  if (code != TSDB_CODE_SUCCESS) {
135!
992
    qError("%s failed at line %d since %s. task:%s", __func__, lino, tstrerror(code), GET_TASKID(pTaskInfo));
×
993
  }
994
  return code;
135✔
995
}
996

997
void resetStreamFillSup(SStreamFillSupporter* pFillSup) {
5,812✔
998
  _hash_fn_t hashFn = taosGetDefaultHashFunction(TSDB_DATA_TYPE_BINARY);
5,812✔
999
  SSHashObj* pNewMap = tSimpleHashInit(16, hashFn);
5,812✔
1000
  if (pNewMap != NULL) {
5,812!
1001
    tSimpleHashCleanup(pFillSup->pResMap);
5,812✔
1002
    pFillSup->pResMap = pNewMap;
5,812✔
1003
  } else {
1004
    tSimpleHashClear(pFillSup->pResMap);
×
1005
  }
1006
  pFillSup->hasDelete = false;
5,812✔
1007
}
5,812✔
1008
void resetStreamFillInfo(SStreamFillOperatorInfo* pInfo) {
3,147✔
1009
  resetStreamFillSup(pInfo->pFillSup);
3,147✔
1010
  taosArrayClear(pInfo->pFillInfo->delRanges);
3,147✔
1011
  pInfo->pFillInfo->delIndex = 0;
3,147✔
1012
}
3,147✔
1013

1014
int32_t doApplyStreamScalarCalculation(SOperatorInfo* pOperator, SSDataBlock* pSrcBlock,
850✔
1015
                                              SSDataBlock* pDstBlock) {
1016
  int32_t                  code = TSDB_CODE_SUCCESS;
850✔
1017
  int32_t                  lino = 0;
850✔
1018
  SStreamFillOperatorInfo* pInfo = pOperator->info;
850✔
1019
  SExprSupp*               pSup = &pOperator->exprSupp;
850✔
1020
  SExecTaskInfo*           pTaskInfo = pOperator->pTaskInfo;
850✔
1021

1022
  blockDataCleanup(pDstBlock);
850✔
1023
  code = blockDataEnsureCapacity(pDstBlock, pSrcBlock->info.rows);
850✔
1024
  QUERY_CHECK_CODE(code, lino, _end);
850!
1025

1026
  code = setInputDataBlock(pSup, pSrcBlock, TSDB_ORDER_ASC, MAIN_SCAN, false);
850✔
1027
  QUERY_CHECK_CODE(code, lino, _end);
850!
1028
  code = projectApplyFunctions(pSup->pExprInfo, pDstBlock, pSrcBlock, pSup->pCtx, pSup->numOfExprs, NULL);
850✔
1029
  QUERY_CHECK_CODE(code, lino, _end);
850!
1030

1031
  pDstBlock->info.rows = 0;
850✔
1032
  pSup = &pInfo->pFillSup->notFillExprSup;
850✔
1033
  code = setInputDataBlock(pSup, pSrcBlock, TSDB_ORDER_ASC, MAIN_SCAN, false);
850✔
1034
  QUERY_CHECK_CODE(code, lino, _end);
850!
1035
  code = projectApplyFunctions(pSup->pExprInfo, pDstBlock, pSrcBlock, pSup->pCtx, pSup->numOfExprs, NULL);
850✔
1036
  QUERY_CHECK_CODE(code, lino, _end);
850!
1037

1038
  pDstBlock->info.id.groupId = pSrcBlock->info.id.groupId;
850✔
1039

1040
  code = blockDataUpdateTsWindow(pDstBlock, pInfo->primaryTsCol);
850✔
1041

1042
_end:
850✔
1043
  if (code != TSDB_CODE_SUCCESS) {
850!
1044
    qError("%s failed at line %d since %s. task:%s", __func__, lino, tstrerror(code), GET_TASKID(pTaskInfo));
×
1045
  }
1046
  return code;
850✔
1047
}
1048

1049
static int32_t doStreamFillNext(SOperatorInfo* pOperator, SSDataBlock** ppRes) {
3,166✔
1050
  int32_t                  code = TSDB_CODE_SUCCESS;
3,166✔
1051
  int32_t                  lino = 0;
3,166✔
1052
  SStreamFillOperatorInfo* pInfo = pOperator->info;
3,166✔
1053
  SExecTaskInfo*           pTaskInfo = pOperator->pTaskInfo;
3,166✔
1054

1055
  if (pOperator->status == OP_EXEC_DONE) {
3,166!
1056
    (*ppRes) = NULL;
×
1057
    return code;
×
1058
  }
1059
  blockDataCleanup(pInfo->pRes);
3,166✔
1060
  if (hasRemainCalc(pInfo->pFillInfo) ||
3,168✔
1061
      (pInfo->pFillInfo->pos != FILL_POS_INVALID && pInfo->pFillInfo->needFill == true)) {
3,158!
1062
    doStreamFillRange(pInfo->pFillInfo, pInfo->pFillSup, pInfo->pRes);
10✔
1063
    if (pInfo->pRes->info.rows > 0) {
10!
1064
      printDataBlock(pInfo->pRes, getStreamOpName(pOperator->operatorType), GET_TASKID(pTaskInfo));
10✔
1065
      (*ppRes) = pInfo->pRes;
10✔
1066
      return code;
10✔
1067
    }
1068
  }
1069
  if (pOperator->status == OP_RES_TO_RETURN) {
3,158✔
1070
    doDeleteFillFinalize(pOperator);
46✔
1071
    if (pInfo->pRes->info.rows > 0) {
46!
1072
      printDataBlock(pInfo->pRes, getStreamOpName(pOperator->operatorType), GET_TASKID(pTaskInfo));
×
1073
      (*ppRes) = pInfo->pRes;
×
1074
      return code;
×
1075
    }
1076
    setOperatorCompleted(pOperator);
46✔
1077
    resetStreamFillInfo(pInfo);
46✔
1078
    (*ppRes) = NULL;
46✔
1079
    return code;
46✔
1080
  }
1081

1082
  SSDataBlock*   fillResult = NULL;
3,112✔
1083
  SOperatorInfo* downstream = pOperator->pDownstream[0];
3,112✔
1084
  while (1) {
1085
    if (pInfo->srcRowIndex >= pInfo->pSrcBlock->info.rows || pInfo->pSrcBlock->info.rows == 0) {
3,157✔
1086
      // If there are delete datablocks, we receive  them first.
1087
      SSDataBlock* pBlock = getNextBlockFromDownstream(pOperator, 0);
3,154✔
1088
      if (pBlock == NULL) {
3,153✔
1089
        pOperator->status = OP_RES_TO_RETURN;
2,147✔
1090
        pInfo->pFillInfo->preRowKey = INT64_MIN;
2,147✔
1091
        if (pInfo->pRes->info.rows > 0) {
2,147!
1092
          printDataBlock(pInfo->pRes, getStreamOpName(pOperator->operatorType), GET_TASKID(pTaskInfo));
×
1093
          (*ppRes) = pInfo->pRes;
×
1094
          return code;
×
1095
        }
1096
        break;
2,147✔
1097
      }
1098
      printSpecDataBlock(pBlock, getStreamOpName(pOperator->operatorType), "recv", GET_TASKID(pTaskInfo));
1,006✔
1099

1100
      if (pInfo->pFillInfo->curGroupId != pBlock->info.id.groupId) {
1,005✔
1101
        pInfo->pFillInfo->curGroupId = pBlock->info.id.groupId;
287✔
1102
        pInfo->pFillInfo->preRowKey = INT64_MIN;
287✔
1103
      }
1104

1105
      pInfo->pFillSup->winRange = pTaskInfo->streamInfo.fillHistoryWindow;
1,005✔
1106
      if (pInfo->pFillSup->winRange.ekey <= 0) {
1,005!
1107
        pInfo->pFillSup->winRange.ekey = INT64_MAX;
×
1108
      }
1109

1110
      switch (pBlock->info.type) {
1,005!
1111
        case STREAM_RETRIEVE:
5✔
1112
          (*ppRes) = pBlock;
5✔
1113
          return code;
5✔
1114
        case STREAM_DELETE_RESULT: {
135✔
1115
          pInfo->pSrcDelBlock = pBlock;
135✔
1116
          pInfo->srcDelRowIndex = 0;
135✔
1117
          blockDataCleanup(pInfo->pDelRes);
135✔
1118
          pInfo->pFillSup->hasDelete = true;
135✔
1119
          code = doDeleteFillResult(pOperator);
135✔
1120
          QUERY_CHECK_CODE(code, lino, _end);
135!
1121

1122
          if (pInfo->pDelRes->info.rows > 0) {
135✔
1123
            printDataBlock(pInfo->pDelRes, getStreamOpName(pOperator->operatorType), GET_TASKID(pTaskInfo));
90✔
1124
            (*ppRes) = pInfo->pDelRes;
90✔
1125
            return code;
90✔
1126
          }
1127
          continue;
45✔
1128
        } break;
1129
        case STREAM_NORMAL:
727✔
1130
        case STREAM_INVALID:
1131
        case STREAM_PULL_DATA: {
1132
          code = doApplyStreamScalarCalculation(pOperator, pBlock, pInfo->pSrcBlock);
727✔
1133
          QUERY_CHECK_CODE(code, lino, _end);
727!
1134

1135
          memcpy(pInfo->pSrcBlock->info.parTbName, pBlock->info.parTbName, TSDB_TABLE_NAME_LEN);
727✔
1136
          pInfo->srcRowIndex = -1;
727✔
1137
        } break;
727✔
1138
        case STREAM_CHECKPOINT:
138✔
1139
        case STREAM_CREATE_CHILD_TABLE: {
1140
          (*ppRes) = pBlock;
138✔
1141
          return code;
138✔
1142
        } break;
1143
        default:
×
1144
          return TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR;
×
1145
      }
1146
    }
1147

1148
    doStreamFillImpl(pOperator);
730✔
1149
    code = doFilter(pInfo->pRes, pOperator->exprSupp.pFilterInfo, &pInfo->matchInfo);
730✔
1150
    QUERY_CHECK_CODE(code, lino, _end);
730!
1151

1152
    memcpy(pInfo->pRes->info.parTbName, pInfo->pSrcBlock->info.parTbName, TSDB_TABLE_NAME_LEN);
730✔
1153
    pOperator->resultInfo.totalRows += pInfo->pRes->info.rows;
730✔
1154
    if (pInfo->pRes->info.rows > 0) {
730!
1155
      break;
730✔
1156
    }
1157
  }
1158
  if (pOperator->status == OP_RES_TO_RETURN) {
2,877✔
1159
    doDeleteFillFinalize(pOperator);
2,147✔
1160
  }
1161

1162
  if (pInfo->pRes->info.rows == 0) {
2,879✔
1163
    setOperatorCompleted(pOperator);
2,103✔
1164
    resetStreamFillInfo(pInfo);
2,103✔
1165
    (*ppRes) = NULL;
2,103✔
1166
    return code;
2,103✔
1167
  }
1168

1169
  pOperator->resultInfo.totalRows += pInfo->pRes->info.rows;
776✔
1170
  printDataBlock(pInfo->pRes, getStreamOpName(pOperator->operatorType), GET_TASKID(pTaskInfo));
776✔
1171
  (*ppRes) = pInfo->pRes;
776✔
1172
  return code;
776✔
1173

1174
_end:
×
1175
  if (code != TSDB_CODE_SUCCESS) {
×
1176
    qError("%s failed at line %d since %s. task:%s", __func__, lino, tstrerror(code), GET_TASKID(pTaskInfo));
×
1177
    pTaskInfo->code = code;
×
1178
    T_LONG_JMP(pTaskInfo->env, code);
×
1179
  }
1180
  setOperatorCompleted(pOperator);
×
1181
  resetStreamFillInfo(pInfo);
×
1182
  (*ppRes) = NULL;
×
1183
  return code;
×
1184
}
1185

1186
static void resetForceFillWindow(SResultRowData* pRowData) {
656✔
1187
  pRowData->key = INT64_MIN;
656✔
1188
  pRowData->pRowVal = NULL;
656✔
1189
}
656✔
1190

1191
void doBuildForceFillResultImpl(SOperatorInfo* pOperator, SStreamFillSupporter* pFillSup,
509✔
1192
                                SStreamFillInfo* pFillInfo, SSDataBlock* pBlock, SGroupResInfo* pGroupResInfo) {
1193
  int32_t code = TSDB_CODE_SUCCESS;
509✔
1194
  int32_t lino = 0;
509✔
1195

1196
  SStreamFillOperatorInfo* pInfo = pOperator->info;
509✔
1197
  bool                     res = false;
509✔
1198
  int32_t                  numOfRows = getNumOfTotalRes(pGroupResInfo);
509✔
1199
  for (; pGroupResInfo->index < numOfRows; pGroupResInfo->index++) {
1,165✔
1200
    SWinKey* pKey = (SWinKey*)taosArrayGet(pGroupResInfo->pRows, pGroupResInfo->index);
806✔
1201
    if (pBlock->info.id.groupId == 0) {
806✔
1202
      pBlock->info.id.groupId = pKey->groupId;
509✔
1203
    } else if (pBlock->info.id.groupId != pKey->groupId) {
297✔
1204
      break;
150✔
1205
    }
1206

1207
    SRowBuffPos* pValPos = NULL;
656✔
1208
    int32_t      len = 0;
656✔
1209
    int32_t      winCode = TSDB_CODE_SUCCESS;
656✔
1210
    code = pInfo->stateStore.streamStateFillGet(pInfo->pState, pKey, (void**)&pValPos, &len, &winCode);
656✔
1211
    QUERY_CHECK_CODE(code, lino, _end);
656!
1212
    qDebug("===stream=== build force fill res. key:%" PRId64 ",groupId:%" PRId64".res:%d", pKey->ts, pKey->groupId, winCode);
656!
1213
    if (winCode == TSDB_CODE_SUCCESS) {
656✔
1214
      pFillSup->cur.key = pKey->ts;
186✔
1215
      pFillSup->cur.pRowVal = pValPos->pRowBuff;
186✔
1216
      code = buildFillResult(&pFillSup->cur, pFillSup, pKey->ts, pBlock, &res, false);
186✔
1217
      QUERY_CHECK_CODE(code, lino, _end);
186!
1218
      resetForceFillWindow(&pFillSup->cur);
186✔
1219
      releaseOutputBuf(pInfo->pState, pValPos, &pInfo->stateStore);
186✔
1220
    } else {
1221
      SWinKey      preKey = {.ts = INT64_MIN, .groupId = pKey->groupId};
470✔
1222
      SRowBuffPos* prePos = NULL;
470✔
1223
      int32_t      preVLen = 0;
470✔
1224
      code = pInfo->stateStore.streamStateFillGetPrev(pInfo->pState, pKey, &preKey,
470✔
1225
                                                      (void**)&prePos, &preVLen, &winCode);
1226
      QUERY_CHECK_CODE(code, lino, _end);
470!
1227
      if (winCode == TSDB_CODE_SUCCESS) {
470!
1228
        pFillSup->cur.key = pKey->ts;
470✔
1229
        pFillSup->cur.pRowVal = prePos->pRowBuff;
470✔
1230
        if (pFillInfo->type == TSDB_FILL_PREV) {
470✔
1231
          code = buildFillResult(&pFillSup->cur, pFillSup, pKey->ts, pBlock, &res, true);
299✔
1232
          QUERY_CHECK_CODE(code, lino, _end);
299!
1233
        } else {
1234
          copyNotFillExpData(pFillSup, pFillInfo);
171✔
1235
          pFillInfo->pResRow->key = pKey->ts;
171✔
1236
          code = buildFillResult(pFillInfo->pResRow, pFillSup, pKey->ts, pBlock, &res, true);
171✔
1237
          QUERY_CHECK_CODE(code, lino, _end);
171!
1238
        }
1239
        resetForceFillWindow(&pFillSup->cur);
470✔
1240
      }
1241
      releaseOutputBuf(pInfo->pState, prePos, &pInfo->stateStore);
470✔
1242
    }
1243
  }
1244

1245
_end:
359✔
1246
  if (code != TSDB_CODE_SUCCESS) {
509!
1247
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
1248
  }
1249
}
509✔
1250

1251
void doBuildForceFillResult(SOperatorInfo* pOperator, SStreamFillSupporter* pFillSup, SStreamFillInfo* pFillInfo,
1,507✔
1252
                            SSDataBlock* pBlock, SGroupResInfo* pGroupResInfo) {
1253
  blockDataCleanup(pBlock);
1,507✔
1254
  if (!hasRemainResults(pGroupResInfo)) {
1,505✔
1255
    return;
998✔
1256
  }
1257

1258
  // clear the existed group id
1259
  pBlock->info.id.groupId = 0;
509✔
1260
  doBuildForceFillResultImpl(pOperator, pFillSup, pFillInfo, pBlock, pGroupResInfo);
509✔
1261
}
1262

1263
static int32_t buildForceFillResult(SOperatorInfo* pOperator, SSDataBlock** ppRes) {
1,507✔
1264
  int32_t                  code = TSDB_CODE_SUCCESS;
1,507✔
1265
  int32_t                  lino = 0;
1,507✔
1266
  SStreamFillOperatorInfo* pInfo = pOperator->info;
1,507✔
1267
  uint16_t                 opType = pOperator->operatorType;
1,507✔
1268
  SExecTaskInfo*           pTaskInfo = pOperator->pTaskInfo;
1,507✔
1269

1270
  doBuildForceFillResult(pOperator, pInfo->pFillSup, pInfo->pFillInfo, pInfo->pRes, &pInfo->groupResInfo);
1,507✔
1271
  if (pInfo->pRes->info.rows != 0) {
1,507✔
1272
    printDataBlock(pInfo->pRes, getStreamOpName(opType), GET_TASKID(pTaskInfo));
509✔
1273
    (*ppRes) = pInfo->pRes;
509✔
1274
    goto _end;
509✔
1275
  }
1276

1277
  (*ppRes) = NULL;
998✔
1278

1279
_end:
1,507✔
1280
  if (code != TSDB_CODE_SUCCESS) {
1,507!
1281
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
1282
  }
1283
  return code;
1,506✔
1284
}
1285

1286
static void keepResultInStateBuf(SStreamFillOperatorInfo* pInfo, uint64_t groupId, SResultRowData* pRow) {
186✔
1287
  int32_t code = TSDB_CODE_SUCCESS;
186✔
1288
  int32_t lino = 0;
186✔
1289

1290
  SWinKey      key = {.groupId = groupId, .ts = pRow->key};
186✔
1291
  int32_t      curVLen = 0;
186✔
1292
  SRowBuffPos* pStatePos = NULL;
186✔
1293
  int32_t      winCode = TSDB_CODE_SUCCESS;
186✔
1294
  code = pInfo->stateStore.streamStateFillAddIfNotExist(pInfo->pState, &key, (void**)&pStatePos,
186✔
1295
                                                        &curVLen, &winCode);
1296
  QUERY_CHECK_CODE(code, lino, _end);
186!
1297
  memcpy(pStatePos->pRowBuff, pRow->pRowVal, pInfo->pFillSup->rowSize);
186✔
1298
  qDebug("===stream===fill operator save key ts:%" PRId64 " group id:%" PRIu64 "  code:%d", key.ts, key.groupId, code);
186!
1299

1300
_end:
×
1301
  if (code != TSDB_CODE_SUCCESS) {
186!
1302
    qError("%s failed at line %d since %s", __func__, __LINE__, tstrerror(code));
×
1303
  }
1304
}
186✔
1305

1306
int32_t keepBlockRowInStateBuf(SStreamFillOperatorInfo* pInfo, SStreamFillInfo* pFillInfo, SSDataBlock* pBlock, TSKEY* tsCol,
186✔
1307
                               int32_t rowId, uint64_t groupId, int32_t rowSize) {
1308
  int32_t code = TSDB_CODE_SUCCESS;
186✔
1309
  int32_t lino = 0;
186✔
1310
  TSKEY ts = tsCol[rowId];
186✔
1311
  pFillInfo->nextRowKey = ts;
186✔
1312
  TAOS_MEMSET(pFillInfo->pTempBuff, 0, rowSize);
186✔
1313
  SResultRowData tmpNextRow = {.key = ts, .pRowVal = pFillInfo->pTempBuff};
186✔
1314

1315
  transBlockToResultRow(pBlock, rowId, ts, &tmpNextRow);
186✔
1316
  keepResultInStateBuf(pInfo, groupId, &tmpNextRow);
186✔
1317

1318
_end:
186✔
1319
  if (code != TSDB_CODE_SUCCESS) {
186!
1320
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
1321
  }
1322
  return code;
186✔
1323
}
1324

1325
// force window close impl
1326
static int32_t doStreamForceFillImpl(SOperatorInfo* pOperator) {
123✔
1327
  int32_t                  code = TSDB_CODE_SUCCESS;
123✔
1328
  int32_t                  lino = 0;
123✔
1329
  SStreamFillOperatorInfo* pInfo = pOperator->info;
123✔
1330
  SExecTaskInfo*           pTaskInfo = pOperator->pTaskInfo;
123✔
1331
  SStreamFillSupporter*    pFillSup = pInfo->pFillSup;
123✔
1332
  SStreamFillInfo*         pFillInfo = pInfo->pFillInfo;
123✔
1333
  SSDataBlock*             pBlock = pInfo->pSrcBlock;
123✔
1334
  uint64_t                 groupId = pBlock->info.id.groupId;
123✔
1335
  SColumnInfoData*         pTsCol = taosArrayGet(pInfo->pSrcBlock->pDataBlock, pInfo->primaryTsCol);
123✔
1336
  TSKEY*                   tsCol = (TSKEY*)pTsCol->pData;
123✔
1337
  for (int32_t i = 0; i < pBlock->info.rows; i++){
309✔
1338
    code = keepBlockRowInStateBuf(pInfo, pFillInfo, pBlock, tsCol, i, groupId, pFillSup->rowSize);
186✔
1339
    QUERY_CHECK_CODE(code, lino, _end);
186!
1340

1341
    int32_t size =  taosArrayGetSize(pInfo->pCloseTs);
186✔
1342
    if (size > 0) {
186!
1343
      TSKEY* pTs = (TSKEY*) taosArrayGet(pInfo->pCloseTs, 0);
186✔
1344
      TSKEY  resTs = tsCol[i];
186✔
1345
      while (resTs < (*pTs)) {
305✔
1346
        SWinKey key = {.groupId = groupId, .ts = resTs};
147✔
1347
        void* pPushRes = taosArrayPush(pInfo->pUpdated, &key);
147✔
1348
        QUERY_CHECK_NULL(pPushRes, code, lino, _end, terrno);
147!
1349

1350
        if (IS_FILL_CONST_VALUE(pFillSup->type)) {
147!
1351
          break;
1352
        }
1353
        resTs = taosTimeAdd(resTs, pFillSup->interval.sliding, pFillSup->interval.slidingUnit,
119✔
1354
                            pFillSup->interval.precision, NULL);
119✔
1355
      }
1356
    }
1357
  }
1358
  code = pInfo->stateStore.streamStateGroupPut(pInfo->pState, groupId, NULL, 0);
123✔
1359
  QUERY_CHECK_CODE(code, lino, _end);
123!
1360

1361
_end:
123✔
1362
  if (code != TSDB_CODE_SUCCESS) {
123!
1363
    qError("%s failed at line %d since %s. task:%s", __func__, lino, tstrerror(code), GET_TASKID(pTaskInfo));
×
1364
  }
1365
  return code;
123✔
1366
}
1367

1368
int32_t buildAllResultKey(SStateStore* pStateStore, SStreamState* pState, TSKEY ts, SArray* pUpdated) {
3,345✔
1369
  int32_t          code = TSDB_CODE_SUCCESS;
3,345✔
1370
  int32_t          lino = 0;
3,345✔
1371
  int64_t          groupId = 0;
3,345✔
1372
  SStreamStateCur* pCur = pStateStore->streamStateGroupGetCur(pState);
3,345✔
1373
  while (1) {  
1,444✔
1374
    int32_t winCode = pStateStore->streamStateGroupGetKVByCur(pCur, &groupId, NULL, NULL);
4,788✔
1375
    if (winCode != TSDB_CODE_SUCCESS) {
4,789✔
1376
      break;
3,345✔
1377
    }
1378
    SWinKey key = {.ts = ts, .groupId = groupId};
1,444✔
1379
    void* pPushRes = taosArrayPush(pUpdated, &key);
1,444✔
1380
    QUERY_CHECK_NULL(pPushRes, code, lino, _end, terrno);
1,444!
1381

1382
    pStateStore->streamStateGroupCurNext(pCur);
1,444✔
1383
  }
1384
  pStateStore->streamStateFreeCur(pCur);
3,345✔
1385
  pCur = NULL;
3,345✔
1386

1387
_end:
3,345✔
1388
  if (code != TSDB_CODE_SUCCESS) {
3,345!
1389
    pStateStore->streamStateFreeCur(pCur);
×
1390
    pCur = NULL;
×
1391
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
1392
  }
1393
  return code;
3,345✔
1394
}
1395

1396
void removeDuplicateResult(SArray* pTsArrray, __compar_fn_t fn) {
998✔
1397
  taosArraySort(pTsArrray, fn);
998✔
1398
  taosArrayRemoveDuplicate(pTsArrray, fn, NULL);
998✔
1399
}
998✔
1400

1401
// force window close
1402
static int32_t doStreamForceFillNext(SOperatorInfo* pOperator, SSDataBlock** ppRes) {
1,512✔
1403
  int32_t                  code = TSDB_CODE_SUCCESS;
1,512✔
1404
  int32_t                  lino = 0;
1,512✔
1405
  SStreamFillOperatorInfo* pInfo = pOperator->info;
1,512✔
1406
  SExecTaskInfo*           pTaskInfo = pOperator->pTaskInfo;
1,512✔
1407

1408
  if (pOperator->status == OP_EXEC_DONE) {
1,512!
1409
    (*ppRes) = NULL;
×
1410
    return code;
×
1411
  }
1412

1413
  if (pOperator->status == OP_RES_TO_RETURN) {
1,512✔
1414
    SSDataBlock* resBlock = NULL;
509✔
1415
    code = buildForceFillResult(pOperator, &resBlock);
509✔
1416
    QUERY_CHECK_CODE(code, lino, _end);
509!
1417

1418
    if (resBlock != NULL) {
509✔
1419
      (*ppRes) = resBlock;
150✔
1420
      goto _end;
150✔
1421
    }
1422

1423
    pInfo->stateStore.streamStateClearExpiredState(pInfo->pState, 1, INT64_MAX);
359✔
1424
    resetStreamFillInfo(pInfo);
359✔
1425
    setStreamOperatorCompleted(pOperator);
359✔
1426
    (*ppRes) = NULL;
359✔
1427
    goto _end;
359✔
1428
  }
1429

1430
  SSDataBlock*   fillResult = NULL;
1,003✔
1431
  SOperatorInfo* downstream = pOperator->pDownstream[0];
1,003✔
1432
  while (1) {
1,119✔
1433
    SSDataBlock* pBlock = getNextBlockFromDownstream(pOperator, 0);
2,122✔
1434
    if (pBlock == NULL) {
2,122✔
1435
      pOperator->status = OP_RES_TO_RETURN;
998✔
1436
      qDebug("===stream===return data:%s.", getStreamOpName(pOperator->operatorType));
998!
1437
      break;
998✔
1438
    }
1439
    printSpecDataBlock(pBlock, getStreamOpName(pOperator->operatorType), "recv", GET_TASKID(pTaskInfo));
1,124✔
1440
    setStreamOperatorState(&pInfo->basic, pBlock->info.type);
1,124✔
1441

1442
    switch (pBlock->info.type) {
1,124!
1443
      case STREAM_NORMAL:
123✔
1444
      case STREAM_INVALID: {
1445
        code = doApplyStreamScalarCalculation(pOperator, pBlock, pInfo->pSrcBlock);
123✔
1446
        QUERY_CHECK_CODE(code, lino, _end);
123!
1447

1448
        memcpy(pInfo->pSrcBlock->info.parTbName, pBlock->info.parTbName, TSDB_TABLE_NAME_LEN);
123✔
1449
        pInfo->srcRowIndex = -1;
123✔
1450
      } break;
123✔
1451
      case STREAM_CHECKPOINT: {
2✔
1452
        pInfo->stateStore.streamStateCommit(pInfo->pState);
2✔
1453
        (*ppRes) = pBlock;
2✔
1454
        goto _end;
2✔
1455
      } break;
1456
      case STREAM_CREATE_CHILD_TABLE: {
3✔
1457
        (*ppRes) = pBlock;
3✔
1458
        goto _end;
3✔
1459
      } break;
1460
      case STREAM_GET_RESULT: {
1,992✔
1461
        void* pPushRes = taosArrayPush(pInfo->pCloseTs, &pBlock->info.window.skey);
996✔
1462
        QUERY_CHECK_NULL(pPushRes, code, lino, _end, terrno);
996!
1463
        continue;
996✔
1464
      }
1465
      default:
×
1466
        code = TSDB_CODE_QRY_EXECUTOR_INTERNAL_ERROR;
×
1467
        QUERY_CHECK_CODE(code, lino, _end);
×
1468
    }
1469

1470
    code = doStreamForceFillImpl(pOperator);
123✔
1471
    QUERY_CHECK_CODE(code, lino, _end);
123!
1472
  }
1473

1474
  for (int32_t i = 0; i < taosArrayGetSize(pInfo->pCloseTs); i++) {
1,994✔
1475
    TSKEY ts = *(TSKEY*) taosArrayGet(pInfo->pCloseTs, i);
996✔
1476
    code = buildAllResultKey(&pInfo->stateStore, pInfo->pState, ts, pInfo->pUpdated);
996✔
1477
    QUERY_CHECK_CODE(code, lino, _end);
996!
1478
  }
1479
  taosArrayClear(pInfo->pCloseTs);
998✔
1480
  removeDuplicateResult(pInfo->pUpdated, winKeyCmprImpl);
998✔
1481

1482
  initMultiResInfoFromArrayList(&pInfo->groupResInfo, pInfo->pUpdated);
998✔
1483
  pInfo->groupResInfo.freeItem = false;
998✔
1484

1485
  pInfo->pUpdated = taosArrayInit(1024, sizeof(SWinKey));
998✔
1486
  QUERY_CHECK_NULL(pInfo->pUpdated, code, lino, _end, terrno);
998!
1487

1488
  code = blockDataEnsureCapacity(pInfo->pRes, pOperator->resultInfo.capacity);
998✔
1489
  QUERY_CHECK_CODE(code, lino, _end);
998!
1490

1491
  code = buildForceFillResult(pOperator, ppRes);
998✔
1492
  QUERY_CHECK_CODE(code, lino, _end);
997!
1493

1494
  if ((*ppRes) == NULL) {
997✔
1495
    pInfo->stateStore.streamStateClearExpiredState(pInfo->pState, 1, INT64_MAX);
639✔
1496
    resetStreamFillInfo(pInfo);
639✔
1497
    setStreamOperatorCompleted(pOperator);
639✔
1498
  }
1499

1500
_end:
358✔
1501
  if (code != TSDB_CODE_SUCCESS) {
1,512!
1502
    qError("%s failed at line %d since %s. task:%s", __func__, lino, tstrerror(code), GET_TASKID(pTaskInfo));
×
1503
    pTaskInfo->code = code;
×
1504
  }
1505
  return code;
1,512✔
1506
}
1507

1508
static int32_t initResultBuf(SSDataBlock* pInputRes, SStreamFillSupporter* pFillSup) {
451✔
1509
  int32_t numOfCols = taosArrayGetSize(pInputRes->pDataBlock);
451✔
1510
  pFillSup->rowSize = sizeof(SResultCellData) * numOfCols;
451✔
1511
  for (int i = 0; i < numOfCols; i++) {
7,724✔
1512
    SColumnInfoData* pCol = taosArrayGet(pInputRes->pDataBlock, i);
7,273✔
1513
    pFillSup->rowSize += pCol->info.bytes;
7,273✔
1514
  }
1515
  pFillSup->next.key = INT64_MIN;
451✔
1516
  pFillSup->nextNext.key = INT64_MIN;
451✔
1517
  pFillSup->prev.key = INT64_MIN;
451✔
1518
  pFillSup->cur.key = INT64_MIN;
451✔
1519
  pFillSup->next.pRowVal = NULL;
451✔
1520
  pFillSup->nextNext.pRowVal = NULL;
451✔
1521
  pFillSup->prev.pRowVal = NULL;
451✔
1522
  pFillSup->cur.pRowVal = NULL;
451✔
1523

1524
  return TSDB_CODE_SUCCESS;
451✔
1525
}
1526

1527
static SStreamFillSupporter* initStreamFillSup(SStreamFillPhysiNode* pPhyFillNode, SInterval* pInterval,
451✔
1528
                                               SExprInfo* pFillExprInfo, int32_t numOfFillCols, SStorageAPI* pAPI, SSDataBlock* pInputRes) {
1529
  int32_t               code = TSDB_CODE_SUCCESS;
451✔
1530
  int32_t               lino = 0;
451✔
1531
  SStreamFillSupporter* pFillSup = taosMemoryCalloc(1, sizeof(SStreamFillSupporter));
451!
1532
  if (!pFillSup) {
451!
1533
    code = terrno;
×
1534
    QUERY_CHECK_CODE(code, lino, _end);
×
1535
  }
1536
  pFillSup->numOfFillCols = numOfFillCols;
451✔
1537
  int32_t    numOfNotFillCols = 0;
451✔
1538
  SExprInfo* noFillExprInfo = NULL;
451✔
1539

1540
  code = createExprInfo(pPhyFillNode->pNotFillExprs, NULL, &noFillExprInfo, &numOfNotFillCols);
451✔
1541
  QUERY_CHECK_CODE(code, lino, _end);
451!
1542

1543
  pFillSup->pAllColInfo = createFillColInfo(pFillExprInfo, pFillSup->numOfFillCols, noFillExprInfo, numOfNotFillCols,
902✔
1544
                                            NULL, 0, (const SNodeListNode*)(pPhyFillNode->pValues));
451✔
1545
  if (pFillSup->pAllColInfo == NULL) {
451!
1546
    code = terrno;
×
1547
    lino = __LINE__;
×
1548
    destroyExprInfo(noFillExprInfo, numOfNotFillCols);
×
1549
    goto _end;
×
1550
  }
1551

1552
  pFillSup->type = convertFillType(pPhyFillNode->mode);
451✔
1553
  pFillSup->numOfAllCols = pFillSup->numOfFillCols + numOfNotFillCols;
451✔
1554
  pFillSup->interval = *pInterval;
451✔
1555
  pFillSup->pAPI = pAPI;
451✔
1556

1557
  code = initResultBuf(pInputRes, pFillSup);
451✔
1558
  QUERY_CHECK_CODE(code, lino, _end);
451!
1559

1560
  SExprInfo* noFillExpr = NULL;
451✔
1561
  code = createExprInfo(pPhyFillNode->pNotFillExprs, NULL, &noFillExpr, &numOfNotFillCols);
451✔
1562
  QUERY_CHECK_CODE(code, lino, _end);
451!
1563

1564
  code = initExprSupp(&pFillSup->notFillExprSup, noFillExpr, numOfNotFillCols, &pAPI->functionStore);
451✔
1565
  QUERY_CHECK_CODE(code, lino, _end);
451!
1566

1567
  _hash_fn_t hashFn = taosGetDefaultHashFunction(TSDB_DATA_TYPE_BINARY);
451✔
1568
  pFillSup->pResMap = tSimpleHashInit(16, hashFn);
451✔
1569
  QUERY_CHECK_NULL(pFillSup->pResMap, code, lino, _end, terrno);
451!
1570
  pFillSup->hasDelete = false;
451✔
1571
  pFillSup->normalFill = true;
451✔
1572
  pFillSup->pResultRange = taosArrayInit(2, POINTER_BYTES);
451✔
1573

1574

1575
_end:
451✔
1576
  if (code != TSDB_CODE_SUCCESS) {
451!
1577
    destroyStreamFillSupporter(pFillSup);
×
1578
    pFillSup = NULL;
×
1579
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
1580
  }
1581
  return pFillSup;
451✔
1582
}
1583

1584
SStreamFillInfo* initStreamFillInfo(SStreamFillSupporter* pFillSup, SSDataBlock* pRes) {
693✔
1585
  int32_t          code = TSDB_CODE_SUCCESS;
693✔
1586
  int32_t          lino = 0;
693✔
1587
  SStreamFillInfo* pFillInfo = taosMemoryCalloc(1, sizeof(SStreamFillInfo));
693!
1588
  if (!pFillInfo) {
693!
1589
    code = terrno;
×
1590
    QUERY_CHECK_CODE(code, lino, _end);
×
1591
  }
1592

1593
  pFillInfo->start = INT64_MIN;
693✔
1594
  pFillInfo->current = INT64_MIN;
693✔
1595
  pFillInfo->end = INT64_MIN;
693✔
1596
  pFillInfo->preRowKey = INT64_MIN;
693✔
1597
  pFillInfo->needFill = false;
693✔
1598
  pFillInfo->pLinearInfo = taosMemoryCalloc(1, sizeof(SStreamFillLinearInfo));
693!
1599
  if (!pFillInfo) {
693!
1600
    code = terrno;
×
1601
    QUERY_CHECK_CODE(code, lino, _end);
×
1602
  }
1603

1604
  pFillInfo->pLinearInfo->hasNext = false;
693✔
1605
  pFillInfo->pLinearInfo->nextEnd = INT64_MIN;
693✔
1606
  pFillInfo->pLinearInfo->pEndPoints = NULL;
693✔
1607
  pFillInfo->pLinearInfo->pNextEndPoints = NULL;
693✔
1608
  if (pFillSup->type == TSDB_FILL_LINEAR) {
693✔
1609
    pFillInfo->pLinearInfo->pEndPoints = taosArrayInit(pFillSup->numOfAllCols, sizeof(SPoint));
106✔
1610
    if (!pFillInfo->pLinearInfo->pEndPoints) {
106!
1611
      code = terrno;
×
1612
      QUERY_CHECK_CODE(code, lino, _end);
×
1613
    }
1614

1615
    pFillInfo->pLinearInfo->pNextEndPoints = taosArrayInit(pFillSup->numOfAllCols, sizeof(SPoint));
106✔
1616
    if (!pFillInfo->pLinearInfo->pNextEndPoints) {
106!
1617
      code = terrno;
×
1618
      QUERY_CHECK_CODE(code, lino, _end);
×
1619
    }
1620

1621
    for (int32_t i = 0; i < pFillSup->numOfAllCols; i++) {
1,538✔
1622
      SColumnInfoData* pColData = taosArrayGet(pRes->pDataBlock, i);
1,432✔
1623
      if (pColData == NULL) {
1,432✔
1624
        SPoint dummy = {0};
16✔
1625
        dummy.val = taosMemoryCalloc(1, 1);
16!
1626
        void* tmpRes = taosArrayPush(pFillInfo->pLinearInfo->pEndPoints, &dummy);
16✔
1627
        QUERY_CHECK_NULL(tmpRes, code, lino, _end, terrno);
16!
1628

1629
        dummy.val = taosMemoryCalloc(1, 1);
16!
1630
        tmpRes = taosArrayPush(pFillInfo->pLinearInfo->pNextEndPoints, &dummy);
16✔
1631
        QUERY_CHECK_NULL(tmpRes, code, lino, _end, terrno);
16!
1632

1633
        continue;
16✔
1634
      }
1635
      SPoint value = {0};
1,416✔
1636
      value.val = taosMemoryCalloc(1, pColData->info.bytes);
1,416!
1637
      QUERY_CHECK_NULL(value.val, code, lino, _end, terrno);
1,416!
1638

1639
      void* tmpRes = taosArrayPush(pFillInfo->pLinearInfo->pEndPoints, &value);
1,416✔
1640
      QUERY_CHECK_NULL(tmpRes, code, lino, _end, terrno);
1,416!
1641

1642
      value.val = taosMemoryCalloc(1, pColData->info.bytes);
1,416!
1643
      QUERY_CHECK_NULL(value.val, code, lino, _end, terrno);
1,416!
1644

1645
      tmpRes = taosArrayPush(pFillInfo->pLinearInfo->pNextEndPoints, &value);
1,416✔
1646
      QUERY_CHECK_NULL(tmpRes, code, lino, _end, terrno);
1,416!
1647
    }
1648
  }
1649
  pFillInfo->pLinearInfo->winIndex = 0;
693✔
1650

1651
  pFillInfo->pNonFillRow = NULL;
693✔
1652
  pFillInfo->pResRow = NULL;
693✔
1653
  if (pFillSup->type == TSDB_FILL_SET_VALUE || pFillSup->type == TSDB_FILL_SET_VALUE_F ||
693✔
1654
      pFillSup->type == TSDB_FILL_NULL || pFillSup->type == TSDB_FILL_NULL_F) {
554✔
1655
    pFillInfo->pResRow = taosMemoryCalloc(1, sizeof(SResultRowData));
305!
1656
    QUERY_CHECK_NULL(pFillInfo->pResRow, code, lino, _end, terrno);
305!
1657

1658
    pFillInfo->pResRow->key = INT64_MIN;
305✔
1659
    pFillInfo->pResRow->pRowVal = taosMemoryCalloc(1, pFillSup->rowSize);
305!
1660
    QUERY_CHECK_NULL(pFillInfo->pResRow->pRowVal, code, lino, _end, terrno);
305!
1661

1662
    for (int32_t i = 0; i < pFillSup->numOfAllCols; ++i) {
4,039✔
1663
      SColumnInfoData* pColData = taosArrayGet(pRes->pDataBlock, i);
3,734✔
1664
      SResultCellData* pCell = getResultCell(pFillInfo->pResRow, i);
3,734✔
1665
      if (pColData == NULL) {
3,734✔
1666
        pCell->bytes = 1;
73✔
1667
        pCell->type = 4;
73✔
1668
        continue;
73✔
1669
      }
1670
      pCell->bytes = pColData->info.bytes;
3,661✔
1671
      pCell->type = pColData->info.type;
3,661✔
1672
    }
1673

1674
    int32_t numOfResCol = taosArrayGetSize(pRes->pDataBlock);
305✔
1675
    if (numOfResCol < pFillSup->numOfAllCols) {
305✔
1676
      int32_t* pTmpBuf = (int32_t*)taosMemoryRealloc(pFillSup->pOffsetInfo, pFillSup->numOfAllCols * sizeof(int32_t));
69!
1677
      QUERY_CHECK_NULL(pTmpBuf, code, lino, _end, terrno);
69!
1678
      pFillSup->pOffsetInfo = pTmpBuf;
69✔
1679

1680
      SResultCellData* pCell = getResultCell(pFillInfo->pResRow, numOfResCol - 1);
69✔
1681
      int32_t preLength = pFillSup->pOffsetInfo[numOfResCol - 1] + pCell->bytes + sizeof(SResultCellData);
69✔
1682
      for (int32_t i = numOfResCol; i < pFillSup->numOfAllCols; i++) {
142✔
1683
        pFillSup->pOffsetInfo[i] = preLength;
73✔
1684
        pCell = getResultCell(pFillInfo->pResRow, i);
73✔
1685
        preLength += pCell->bytes + sizeof(SResultCellData);
73✔
1686
      }
1687
    }
1688

1689
    pFillInfo->pNonFillRow = taosMemoryCalloc(1, sizeof(SResultRowData));
305!
1690
    QUERY_CHECK_NULL(pFillInfo->pNonFillRow, code, lino, _end, terrno);
305!
1691
    pFillInfo->pNonFillRow->key = INT64_MIN;
305✔
1692
    pFillInfo->pNonFillRow->pRowVal = taosMemoryCalloc(1, pFillSup->rowSize);
305!
1693
    memcpy(pFillInfo->pNonFillRow->pRowVal, pFillInfo->pResRow->pRowVal, pFillSup->rowSize);
305✔
1694
  }
1695

1696
  pFillInfo->type = pFillSup->type;
693✔
1697
  pFillInfo->delRanges = taosArrayInit(16, sizeof(STimeFillRange));
693✔
1698
  if (!pFillInfo->delRanges) {
693!
1699
    code = terrno;
×
1700
    QUERY_CHECK_CODE(code, lino, _end);
×
1701
  }
1702

1703
  pFillInfo->delIndex = 0;
693✔
1704
  pFillInfo->curGroupId = 0;
693✔
1705
  pFillInfo->hasNext = false;
693✔
1706
  pFillInfo->pTempBuff = taosMemoryCalloc(1, pFillSup->rowSize);
693!
1707
  return pFillInfo;
693✔
1708

1709
_end:
×
1710
  if (code != TSDB_CODE_SUCCESS) {
×
1711
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
1712
  }
1713
  destroyStreamFillInfo(pFillInfo);
×
1714
  return NULL;
×
1715
}
1716

1717
static void setValueForFillInfo(SStreamFillSupporter* pFillSup, SStreamFillInfo* pFillInfo) {
451✔
1718
  if (pFillInfo->type == TSDB_FILL_SET_VALUE || pFillInfo->type == TSDB_FILL_SET_VALUE_F) {
451✔
1719
    for (int32_t i = 0; i < pFillSup->numOfAllCols; ++i) {
1,445✔
1720
      SFillColInfo*    pFillCol = pFillSup->pAllColInfo + i;
1,358✔
1721
      int32_t          slotId = GET_DEST_SLOT_ID(pFillCol);
1,358✔
1722
      SResultCellData* pCell = getResultCell(pFillInfo->pResRow, slotId);
1,358✔
1723
      SVariant*        pVar = &(pFillCol->fillVal);
1,358✔
1724
      if (pCell->type == TSDB_DATA_TYPE_FLOAT) {
1,358!
1725
        float v = 0;
×
1726
        GET_TYPED_DATA(v, float, pVar->nType, &pVar->i, 0);
×
1727
        SET_TYPED_DATA(pCell->pData, pCell->type, v);
×
1728
      } else if (IS_FLOAT_TYPE(pCell->type)) {
1,806!
1729
        double v = 0;
448✔
1730
        GET_TYPED_DATA(v, double, pVar->nType, &pVar->i, 0);
448!
1731
        SET_TYPED_DATA(pCell->pData, pCell->type, v);
448!
1732
      } else if (IS_INTEGER_TYPE(pCell->type)) {
1,733!
1733
        int64_t v = 0;
823✔
1734
        GET_TYPED_DATA(v, int64_t, pVar->nType, &pVar->i, 0);
823!
1735
        SET_TYPED_DATA(pCell->pData, pCell->type, v);
823!
1736
      } else {
1737
        pCell->isNull = true;
87✔
1738
      }
1739
    }
1740
  } else if (pFillInfo->type == TSDB_FILL_NULL || pFillInfo->type == TSDB_FILL_NULL_F) {
364✔
1741
    for (int32_t i = 0; i < pFillSup->numOfAllCols; ++i) {
2,078✔
1742
      SFillColInfo*    pFillCol = pFillSup->pAllColInfo + i;
1,954✔
1743
      int32_t          slotId = GET_DEST_SLOT_ID(pFillCol);
1,954✔
1744
      SResultCellData* pCell = getResultCell(pFillInfo->pResRow, slotId);
1,954✔
1745
      pCell->isNull = true;
1,954✔
1746
    }
1747
  }
1748
}
451✔
1749

1750
int32_t getDownStreamInfo(SOperatorInfo* downstream, int8_t* triggerType, SInterval* pInterval,
451✔
1751
                          int16_t* pOperatorFlag) {
1752
  int32_t code = TSDB_CODE_SUCCESS;
451✔
1753
  int32_t lino = 0;
451✔
1754
  if (IS_NORMAL_INTERVAL_OP(downstream)) {
451✔
1755
    SStreamIntervalOperatorInfo* pInfo = downstream->info;
407✔
1756
    *triggerType = pInfo->twAggSup.calTrigger;
407✔
1757
    *pInterval = pInfo->interval;
407✔
1758
    *pOperatorFlag = pInfo->basic.operatorFlag;
407✔
1759
  } else {
1760
    SStreamIntervalSliceOperatorInfo* pInfo = downstream->info;
44✔
1761
    *triggerType = pInfo->twAggSup.calTrigger;
44✔
1762
    *pInterval = pInfo->interval;
44✔
1763
    pInfo->hasFill = true;
44✔
1764
    *pOperatorFlag = pInfo->basic.operatorFlag;
44✔
1765
  }
1766

1767
  QUERY_CHECK_CODE(code, lino, _end);
451!
1768

1769
_end:
451✔
1770
  if (code != TSDB_CODE_SUCCESS) {
451!
1771
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
1772
  }
1773
  return code;
451✔
1774
}
1775

1776
int32_t initFillOperatorStateBuff(SStreamFillOperatorInfo* pInfo, SStreamState* pState, SStateStore* pStore,
44✔
1777
                                  SReadHandle* pHandle, const char* taskIdStr, SStorageAPI* pApi) {
1778
  int32_t code = TSDB_CODE_SUCCESS;
44✔
1779
  int32_t lino = 0;
44✔
1780

1781
  pInfo->stateStore = *pStore;
44✔
1782
  pInfo->pState = taosMemoryCalloc(1, sizeof(SStreamState));
44!
1783
  QUERY_CHECK_NULL(pInfo->pState, code, lino, _end, terrno);
44!
1784

1785
  *(pInfo->pState) = *pState;
44✔
1786
  pInfo->stateStore.streamStateSetNumber(pInfo->pState, -1, pInfo->primaryTsCol);
44✔
1787
  code = pInfo->stateStore.streamFileStateInit(tsStreamBufferSize, sizeof(SWinKey), pInfo->pFillSup->rowSize, 0, compareTs,
44✔
1788
                                               pInfo->pState, INT64_MAX, taskIdStr, pHandle->checkpointId,
44✔
1789
                                               STREAM_STATE_BUFF_HASH_SORT, &pInfo->pState->pFileState);
44✔
1790
  QUERY_CHECK_CODE(code, lino, _end);
44!
1791

1792
_end:
44✔
1793
  if (code != TSDB_CODE_SUCCESS) {
44!
1794
    qError("%s failed at line %d since %s", __func__, lino, tstrerror(code));
×
1795
  }
1796
  return code;
44✔
1797
}
1798

1799
int32_t createStreamFillOperatorInfo(SOperatorInfo* downstream, SStreamFillPhysiNode* pPhyFillNode,
451✔
1800
                                     SExecTaskInfo* pTaskInfo, SReadHandle* pHandle, SOperatorInfo** pOptrInfo) {
1801
  QRY_PARAM_CHECK(pOptrInfo);
451!
1802

1803
  int32_t                  code = TSDB_CODE_SUCCESS;
451✔
1804
  int32_t                  lino = 0;
451✔
1805
  SStreamFillOperatorInfo* pInfo = taosMemoryCalloc(1, sizeof(SStreamFillOperatorInfo));
451!
1806
  SOperatorInfo*           pOperator = taosMemoryCalloc(1, sizeof(SOperatorInfo));
451!
1807
  if (pInfo == NULL || pOperator == NULL) {
451!
1808
    code = terrno;
×
1809
    QUERY_CHECK_CODE(code, lino, _error);
×
1810
  }
1811

1812
  int32_t    numOfFillCols = 0;
451✔
1813
  SExprInfo* pFillExprInfo = NULL;
451✔
1814

1815
  code = createExprInfo(pPhyFillNode->pFillExprs, NULL, &pFillExprInfo, &numOfFillCols);
451✔
1816
  QUERY_CHECK_CODE(code, lino, _error);
451!
1817

1818
  code = initExprSupp(&pOperator->exprSupp, pFillExprInfo, numOfFillCols, &pTaskInfo->storageAPI.functionStore);
451✔
1819
  QUERY_CHECK_CODE(code, lino, _error);
451!
1820

1821
  pInfo->pSrcBlock = createDataBlockFromDescNode(pPhyFillNode->node.pOutputDataBlockDesc);
451✔
1822
  QUERY_CHECK_NULL(pInfo->pSrcBlock, code, lino, _error, terrno);
451!
1823

1824
  int8_t triggerType = 0;
451✔
1825
  SInterval interval = {0};
451✔
1826
  int16_t opFlag = 0;
451✔
1827
  code = getDownStreamInfo(downstream, &triggerType, &interval, &opFlag);
451✔
1828
  QUERY_CHECK_CODE(code, lino, _error);
451!
1829

1830
  pInfo->pFillSup = initStreamFillSup(pPhyFillNode, &interval, pFillExprInfo, numOfFillCols, &pTaskInfo->storageAPI,
451✔
1831
                                      pInfo->pSrcBlock);
1832
  if (!pInfo->pFillSup) {
451!
1833
    code = TSDB_CODE_FAILED;
×
1834
    QUERY_CHECK_CODE(code, lino, _error);
×
1835
  }
1836

1837
  initResultSizeInfo(&pOperator->resultInfo, 4096);
451✔
1838
  pInfo->pRes = createDataBlockFromDescNode(pPhyFillNode->node.pOutputDataBlockDesc);
451✔
1839
  QUERY_CHECK_NULL(pInfo->pRes, code, lino, _error, terrno);
451!
1840

1841
  code = blockDataEnsureCapacity(pInfo->pRes, pOperator->resultInfo.capacity);
451✔
1842
  QUERY_CHECK_CODE(code, lino, _error);
451!
1843

1844
  code = blockDataEnsureCapacity(pInfo->pSrcBlock, pOperator->resultInfo.capacity);
451✔
1845
  QUERY_CHECK_CODE(code, lino, _error);
451!
1846

1847
  pInfo->pFillInfo = initStreamFillInfo(pInfo->pFillSup, pInfo->pRes);
451✔
1848
  if (!pInfo->pFillInfo) {
451!
1849
    goto _error;
×
1850
  }
1851

1852
  setValueForFillInfo(pInfo->pFillSup, pInfo->pFillInfo);
451✔
1853

1854
  code = createSpecialDataBlock(STREAM_DELETE_RESULT, &pInfo->pDelRes);
451✔
1855
  QUERY_CHECK_CODE(code, lino, _error);
451!
1856

1857
  code = blockDataEnsureCapacity(pInfo->pDelRes, pOperator->resultInfo.capacity);
451✔
1858
  QUERY_CHECK_CODE(code, lino, _error);
451!
1859

1860
  pInfo->pUpdated = taosArrayInit(1024, sizeof(SWinKey));
451✔
1861
  QUERY_CHECK_NULL(pInfo->pUpdated, code, lino, _error, terrno);
451!
1862

1863
  pInfo->pCloseTs = taosArrayInit(1024, sizeof(TSKEY));
451✔
1864
  QUERY_CHECK_NULL(pInfo->pCloseTs, code, lino, _error, terrno);
451!
1865

1866
  pInfo->primaryTsCol = ((STargetNode*)pPhyFillNode->pWStartTs)->slotId;
451✔
1867
  pInfo->primarySrcSlotId = ((SColumnNode*)((STargetNode*)pPhyFillNode->pWStartTs)->pExpr)->slotId;
451✔
1868

1869
  int32_t numOfOutputCols = 0;
451✔
1870
  code = extractColMatchInfo(pPhyFillNode->pFillExprs, pPhyFillNode->node.pOutputDataBlockDesc, &numOfOutputCols,
451✔
1871
                             COL_MATCH_FROM_SLOT_ID, &pInfo->matchInfo);
1872
  QUERY_CHECK_CODE(code, lino, _error);
451!
1873

1874
  code = filterInitFromNode((SNode*)pPhyFillNode->node.pConditions, &pOperator->exprSupp.pFilterInfo, 0);
451✔
1875
  QUERY_CHECK_CODE(code, lino, _error);
451!
1876

1877
  pInfo->srcRowIndex = -1;
451✔
1878
  setOperatorInfo(pOperator, "StreamFillOperator", nodeType(pPhyFillNode), false, OP_NOT_OPENED, pInfo,
451✔
1879
                  pTaskInfo);
1880

1881
  if (triggerType == STREAM_TRIGGER_FORCE_WINDOW_CLOSE) {
451✔
1882
    code = initFillOperatorStateBuff(pInfo, pTaskInfo->streamInfo.pState, &pTaskInfo->storageAPI.stateStore, pHandle,
44✔
1883
                              GET_TASKID(pTaskInfo), &pTaskInfo->storageAPI);
44✔
1884
    QUERY_CHECK_CODE(code, lino, _error);
44!
1885
    pOperator->fpSet = createOperatorFpSet(optrDummyOpenFn, doStreamForceFillNext, NULL, destroyStreamFillOperatorInfo,
44✔
1886
                                           optrDefaultBufFn, NULL, optrDefaultGetNextExtFn, NULL);
1887
  } else if (triggerType == STREAM_TRIGGER_CONTINUOUS_WINDOW_CLOSE) {
407!
1888
    code = initFillOperatorStateBuff(pInfo, pTaskInfo->streamInfo.pState, &pTaskInfo->storageAPI.stateStore, pHandle,
×
1889
                              GET_TASKID(pTaskInfo), &pTaskInfo->storageAPI);
×
1890
    QUERY_CHECK_CODE(code, lino, _error);
×
1891

1892
    initNonBlockAggSupptor(&pInfo->nbSup, &pInfo->pFillSup->interval, downstream);
×
1893
    code = initStreamBasicInfo(&pInfo->basic, pOperator);
×
1894
    QUERY_CHECK_CODE(code, lino, _error);
×
1895

1896
    code = streamClientCheckCfg(&pInfo->nbSup.recParam);
×
1897
    QUERY_CHECK_CODE(code, lino, _error);
×
1898

1899
    pInfo->basic.operatorFlag = opFlag;
×
1900
    if (isFinalOperator(&pInfo->basic)) {
×
1901
      pInfo->nbSup.numOfKeep++;
×
1902
    }
1903
    code = initFillSupRowInfo(pInfo->pFillSup, pInfo->pRes);
×
1904
    QUERY_CHECK_CODE(code, lino, _error);
×
1905
    pOperator->fpSet = createOperatorFpSet(optrDummyOpenFn, doStreamNonblockFillNext, NULL, destroyStreamNonblockFillOperatorInfo,
×
1906
                                           optrDefaultBufFn, NULL, optrDefaultGetNextExtFn, NULL);
1907
  } else {
1908
    pInfo->pState = NULL;
407✔
1909
    pOperator->fpSet = createOperatorFpSet(optrDummyOpenFn, doStreamFillNext, NULL, destroyStreamFillOperatorInfo,
407✔
1910
                                           optrDefaultBufFn, NULL, optrDefaultGetNextExtFn, NULL);
1911
  }
1912
  setOperatorStreamStateFn(pOperator, streamOpReleaseState, streamOpReloadState);
451✔
1913

1914
  code = appendDownstream(pOperator, &downstream, 1);
451✔
1915
  QUERY_CHECK_CODE(code, lino, _error);
451!
1916

1917
  *pOptrInfo = pOperator;
451✔
1918
  return TSDB_CODE_SUCCESS;
451✔
1919

1920
_error:
×
1921
  qError("%s failed at line %d since %s. task:%s", __func__, lino, tstrerror(code), GET_TASKID(pTaskInfo));
×
1922

1923
  if (pInfo != NULL) destroyStreamFillOperatorInfo(pInfo);
×
1924
  destroyOperatorAndDownstreams(pOperator, &downstream, 1);
×
1925
  pTaskInfo->code = code;
×
1926
  return code;
×
1927
}
STATUS · Troubleshooting · Open an Issue · Sales · Support · CAREERS · ENTERPRISE · START FREE · SCHEDULE DEMO
ANNOUNCEMENTS · TWITTER · TOS & SLA · Supported CI Services · What's a CI service? · Automated Testing

© 2026 Coveralls, Inc