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

16
#include "os.h"
17
#include "streamState.h"
L
Liu Jicong 已提交
18
#include "tdatablock.h"
L
Liu Jicong 已提交
19
#include "tdbInt.h"
L
Liu Jicong 已提交
20 21
#include "tmsg.h"
#include "tmsgcb.h"
L
Liu Jicong 已提交
22
#include "tqueue.h"
L
Liu Jicong 已提交
23 24 25 26 27

#ifdef __cplusplus
extern "C" {
#endif

L
Liu Jicong 已提交
28 29
#ifndef _STREAM_H_
#define _STREAM_H_
L
Liu Jicong 已提交
30

L
Liu Jicong 已提交
31 32
typedef struct SStreamTask SStreamTask;

Y
yihaoDeng 已提交
33
#define SSTREAM_TASK_VER 1
34 35
enum {
  STREAM_STATUS__NORMAL = 0,
L
Liu Jicong 已提交
36
  STREAM_STATUS__STOP,
L
Liu Jicong 已提交
37
  STREAM_STATUS__INIT,
L
Liu Jicong 已提交
38
  STREAM_STATUS__FAILED,
39
  STREAM_STATUS__RECOVER,
5
54liuyao 已提交
40
  STREAM_STATUS__PAUSE,
41 42
};

L
Liu Jicong 已提交
43
enum {
L
Liu Jicong 已提交
44 45
  TASK_STATUS__NORMAL = 0,
  TASK_STATUS__DROPPING,
L
Liu Jicong 已提交
46 47
  TASK_STATUS__FAIL,
  TASK_STATUS__STOP,
Y
yihaoDeng 已提交
48
  TASK_STATUS__SCAN_HISTORY,  // stream task scan history data by using tsdbread in the stream scanner
49
  TASK_STATUS__HALT,          // pause, but not be manipulated by user command
50
  TASK_STATUS__PAUSE,         // pause
L
Liu Jicong 已提交
51 52 53
};

enum {
L
Liu Jicong 已提交
54 55 56 57
  TASK_SCHED_STATUS__INACTIVE = 1,
  TASK_SCHED_STATUS__WAITING,
  TASK_SCHED_STATUS__ACTIVE,
  TASK_SCHED_STATUS__FAILED,
L
Liu Jicong 已提交
58
  TASK_SCHED_STATUS__DROPPING,
L
Liu Jicong 已提交
59 60 61 62 63 64 65
};

enum {
  TASK_INPUT_STATUS__NORMAL = 1,
  TASK_INPUT_STATUS__BLOCKED,
  TASK_INPUT_STATUS__RECOVER,
  TASK_INPUT_STATUS__STOP,
L
Liu Jicong 已提交
66
  TASK_INPUT_STATUS__FAILED,
L
Liu Jicong 已提交
67 68 69 70 71 72
};

enum {
  TASK_OUTPUT_STATUS__NORMAL = 1,
  TASK_OUTPUT_STATUS__WAIT,
  TASK_OUTPUT_STATUS__BLOCKED,
L
Liu Jicong 已提交
73 74
};

75 76 77 78 79
enum {
  TASK_TRIGGER_STATUS__INACTIVE = 1,
  TASK_TRIGGER_STATUS__ACTIVE,
};

80
typedef enum {
81 82 83
  TASK_LEVEL__SOURCE = 1,
  TASK_LEVEL__AGG,
  TASK_LEVEL__SINK,
84
} ETASK_LEVEL;
85 86 87 88 89 90 91 92 93

enum {
  TASK_OUTPUT__FIXED_DISPATCH = 1,
  TASK_OUTPUT__SHUFFLE_DISPATCH,
  TASK_OUTPUT__TABLE,
  TASK_OUTPUT__SMA,
  TASK_OUTPUT__FETCH,
};

L
Liu Jicong 已提交
94 95 96 97 98 99
enum {
  STREAM_QUEUE__SUCESS = 1,
  STREAM_QUEUE__FAILED,
  STREAM_QUEUE__PROCESSING,
};

L
Liu Jicong 已提交
100 101 102 103
typedef struct {
  int8_t type;
} SStreamQueueItem;

104 105
typedef void    FTbSink(SStreamTask* pTask, void* vnode, int64_t ver, void* data);
typedef int32_t FTaskExpand(void* ahandle, SStreamTask* pTask, int64_t ver);
L
Liu Jicong 已提交
106 107

typedef struct {
L
Liu Jicong 已提交
108 109 110 111
  int8_t      type;
  int64_t     ver;
  int32_t*    dataRef;
  SPackedData submit;
112
} SStreamDataSubmit;
L
Liu Jicong 已提交
113 114 115 116 117 118

typedef struct {
  int8_t  type;
  int64_t ver;
  SArray* dataRefs;  // SArray<int32_t*>
  SArray* submits;   // SArray<SPackedSubmit>
119
} SStreamMergedSubmit;
120

L
Liu Jicong 已提交
121 122 123
typedef struct {
  int8_t type;

L
Liu Jicong 已提交
124
  int32_t srcVgId;
125
  int32_t srcTaskId;
L
Liu Jicong 已提交
126
  int32_t childId;
L
Liu Jicong 已提交
127
  int64_t sourceVer;
L
Liu Jicong 已提交
128
  int64_t reqId;
L
Liu Jicong 已提交
129

L
Liu Jicong 已提交
130
  SArray* blocks;  // SArray<SSDataBlock>
L
Liu Jicong 已提交
131 132
} SStreamDataBlock;

L
Liu Jicong 已提交
133 134 135 136 137 138
// ref data block, for delete
typedef struct {
  int8_t       type;
  SSDataBlock* pBlock;
} SStreamRefDataBlock;

L
Liu Jicong 已提交
139 140 141 142
typedef struct {
  int8_t type;
} SStreamCheckpoint;

143 144 145 146 147
typedef struct {
  int8_t       type;
  SSDataBlock* pBlock;
} SStreamTrigger;

L
Liu Jicong 已提交
148 149 150 151 152 153 154 155 156 157 158 159 160 161
typedef struct SStreamQueueNode SStreamQueueNode;

struct SStreamQueueNode {
  SStreamQueueItem* item;
  SStreamQueueNode* next;
};

typedef struct {
  SStreamQueueNode* head;
  int64_t           size;
} SStreamQueueRes;

void streamFreeQitem(SStreamQueueItem* data);

5
54liuyao 已提交
162
#if 0
L
Liu Jicong 已提交
163 164 165 166 167 168
bool              streamQueueResEmpty(const SStreamQueueRes* pRes);
int64_t           streamQueueResSize(const SStreamQueueRes* pRes);
SStreamQueueNode* streamQueueResFront(SStreamQueueRes* pRes);
SStreamQueueNode* streamQueueResPop(SStreamQueueRes* pRes);
void              streamQueueResClear(SStreamQueueRes* pRes);
SStreamQueueRes   streamQueueBuildRes(SStreamQueueNode* pNode);
5
54liuyao 已提交
169
#endif
L
Liu Jicong 已提交
170 171 172 173 174

typedef struct {
  SStreamQueueNode* pHead;
} SStreamQueue1;

5
54liuyao 已提交
175
#if 0
L
Liu Jicong 已提交
176 177 178
bool            streamQueueHasTask(const SStreamQueue1* pQueue);
int32_t         streamQueuePush(SStreamQueue1* pQueue, SStreamQueueItem* pItem);
SStreamQueueRes streamQueueGetRes(SStreamQueue1* pQueue);
5
54liuyao 已提交
179
#endif
L
Liu Jicong 已提交
180

L
Liu Jicong 已提交
181 182 183 184
typedef struct {
  STaosQueue* queue;
  STaosQall*  qall;
  void*       qItem;
L
Liu Jicong 已提交
185 186
  int8_t      status;
} SStreamQueue;
L
Liu Jicong 已提交
187

188 189 190
int32_t streamInit();
void    streamCleanUp();

dengyihao's avatar
dengyihao 已提交
191
SStreamQueue* streamQueueOpen(int64_t cap);
L
Liu Jicong 已提交
192 193 194
void          streamQueueClose(SStreamQueue* queue);

static FORCE_INLINE void streamQueueProcessSuccess(SStreamQueue* queue) {
195
  ASSERT(atomic_load_8(&queue->status) == STREAM_QUEUE__PROCESSING);
L
Liu Jicong 已提交
196 197
  queue->qItem = NULL;
  atomic_store_8(&queue->status, STREAM_QUEUE__SUCESS);
L
Liu Jicong 已提交
198 199
}

L
Liu Jicong 已提交
200
static FORCE_INLINE void streamQueueProcessFail(SStreamQueue* queue) {
201
  ASSERT(atomic_load_8(&queue->status) == STREAM_QUEUE__PROCESSING);
L
Liu Jicong 已提交
202 203 204
  atomic_store_8(&queue->status, STREAM_QUEUE__FAILED);
}

205
void* streamQueueNextItem(SStreamQueue* pQueue);
L
Liu Jicong 已提交
206

207
SStreamDataSubmit* streamDataSubmitNew(SPackedData* pData, int32_t type);
208
void               streamDataSubmitDestroy(SStreamDataSubmit* pDataSubmit);
L
Liu Jicong 已提交
209

L
Liu Jicong 已提交
210
typedef struct {
211 212 213
  char*              qmsg;
  void*              pExecutor;   // not applicable to encoder and decoder
  struct SWalReader* pWalReader;  // not applicable to encoder and decoder
L
Liu Jicong 已提交
214 215 216
} STaskExec;

typedef struct {
L
Liu Jicong 已提交
217
  int32_t taskId;
L
Liu Jicong 已提交
218 219 220 221 222
  int32_t nodeId;
  SEpSet  epSet;
} STaskDispatcherFixedEp;

typedef struct {
L
Liu Jicong 已提交
223
  char      stbFullName[TSDB_TABLE_FNAME_LEN];
L
Liu Jicong 已提交
224
  int32_t   waitingRspCnt;
L
Liu Jicong 已提交
225
  SUseDbRsp dbInfo;
L
Liu Jicong 已提交
226 227 228
} STaskDispatcherShuffle;

typedef struct {
L
Liu Jicong 已提交
229
  int64_t         stbUid;
L
Liu Jicong 已提交
230
  char            stbFullName[TSDB_TABLE_FNAME_LEN];
L
Liu Jicong 已提交
231
  SSchemaWrapper* pSchemaWrapper;
232 233 234
  void*           vnode;  // not available to encoder and decoder
  FTbSink*        tbSinkFunc;
  STSchema*       pTSchema;
L
liuyao 已提交
235
  SSHashObj*      pTblInfo;
L
Liu Jicong 已提交
236 237
} STaskSinkTb;

L
Liu Jicong 已提交
238
typedef void FSmaSink(void* vnode, int64_t smaId, const SArray* data);
L
Liu Jicong 已提交
239

L
Liu Jicong 已提交
240
typedef struct {
L
Liu Jicong 已提交
241 242
  int64_t smaId;
  // following are not applicable to encoder and decoder
L
Liu Jicong 已提交
243
  void*     vnode;
L
Liu Jicong 已提交
244
  FSmaSink* smaSink;
L
Liu Jicong 已提交
245 246 247 248 249 250
} STaskSinkSma;

typedef struct {
  int8_t reserved;
} STaskSinkFetch;

251
typedef struct SStreamChildEpInfo {
L
Liu Jicong 已提交
252 253 254
  int32_t nodeId;
  int32_t childId;
  int32_t taskId;
255
  int8_t  dataAllowed;
L
Liu Jicong 已提交
256
  SEpSet  epSet;
L
Liu Jicong 已提交
257 258
} SStreamChildEpInfo;

259 260 261 262 263
typedef struct SStreamId {
  int64_t     streamId;
  int32_t     taskId;
  const char* idStr;
} SStreamId;
L
Liu Jicong 已提交
264

265 266
typedef struct SCheckpointInfo {
  int64_t id;
267 268
  int64_t version;     // offset in WAL
  int64_t currentVer;  // current offset in WAL, not serialize it
269 270
} SCheckpointInfo;

271
typedef struct SStreamStatus {
Y
yihaoDeng 已提交
272 273 274 275 276
  int8_t taskStatus;
  int8_t downstreamReady;  // downstream tasks are all ready now, if this flag is set
  int8_t schedStatus;
  int8_t keepTaskStatus;
  bool   transferState;
277
  bool   appendTranstateBlock;  // has append the transfer state data block already, todo: remove it
Y
yihaoDeng 已提交
278 279
  int8_t timerActive;   // timer is active
  int8_t pauseAllowed;  // allowed task status to be set to be paused
280
} SStreamStatus;
L
Liu Jicong 已提交
281

282
typedef struct SHistDataRange {
283 284
  SVersionRange range;
  STimeWindow   window;
285
} SHistDataRange;
286

287
typedef struct SSTaskBasicInfo {
Y
yihaoDeng 已提交
288 289 290 291 292 293
  int32_t nodeId;  // vgroup id or snode id
  SEpSet  epSet;
  int32_t selfChildId;
  int32_t totalLevel;
  int8_t  taskLevel;
  int8_t  fillHistory;  // is fill history task or not
294
} SSTaskBasicInfo;
295

296
typedef struct SDispatchMsgInfo {
Y
yihaoDeng 已提交
297 298 299 300
  void*   pData;       // current dispatch data
  int16_t msgType;     // dispatch msg type
  int32_t retryCount;  // retry send data count
  int64_t blockingTs;  // output blocking timestamp
301
} SDispatchMsgInfo;
302

303
typedef struct {
304 305 306 307
  int8_t        type;
  int8_t        status;
  SStreamQueue* queue;
} STaskOutputInfo;
308

H
Haojun Liao 已提交
309 310 311 312 313 314
typedef struct {
  int64_t init;
  int64_t step1Start;
  int64_t step2Start;
} STaskTimestamp;

315
struct SStreamTask {
Y
yihaoDeng 已提交
316
  int64_t          ver;
317 318
  SStreamId        id;
  SSTaskBasicInfo  info;
319
  STaskOutputInfo  outputInfo;
320 321 322 323 324 325
  SDispatchMsgInfo msgInfo;
  SStreamStatus    status;
  SCheckpointInfo  chkInfo;
  STaskExec        exec;
  SHistDataRange   dataRange;
  SStreamId        historyTaskId;
H
Haojun Liao 已提交
326
  SStreamId        streamTaskId;
327 328 329
  SArray*          pUpstreamEpInfoList;  // SArray<SStreamChildEpInfo*>, // children info
  int32_t          nextCheckId;
  SArray*          checkpointInfo;  // SArray<SStreamCheckpointInfo>
H
Haojun Liao 已提交
330
  STaskTimestamp   tsInfo;
331
  // output
L
Liu Jicong 已提交
332 333 334
  union {
    STaskDispatcherFixedEp fixedEpDispatcher;
    STaskDispatcherShuffle shuffleDispatcher;
335 336 337
    STaskSinkTb            tbSink;
    STaskSinkSma           smaSink;
    STaskSinkFetch         fetchSink;
L
Liu Jicong 已提交
338 339
  };

340
  int8_t        inputStatus;
L
Liu Jicong 已提交
341
  SStreamQueue* inputQueue;
L
Liu Jicong 已提交
342

343
  // trigger
344 345
  int8_t        triggerStatus;
  int64_t       triggerParam;
346 347
  void*         schedTimer;
  void*         launchTaskTimer;
348 349
  SMsgCb*       pMsgCb;  // msg handle
  SStreamState* pState;  // state backend
350
  SArray*       pRspMsgList;
H
Haojun Liao 已提交
351
  TdThreadMutex lock;
352

353
  // the followings attributes don't be serialized
354
  int32_t             notReadyTasks;
355
  int32_t             numOfWaitingUpstream;
356 357 358 359 360
  int64_t             checkReqId;
  SArray*             checkReqIds;  // shuffle
  int32_t             refCnt;
  int64_t             checkpointingId;
  int32_t             checkpointAlignCnt;
H
Haojun Liao 已提交
361
  int32_t             transferStateAlignCnt;
362
  struct SStreamMeta* pMeta;
L
liuyao 已提交
363
  SSHashObj*          pNameMap;
364
};
L
Liu Jicong 已提交
365

366 367
// meta
typedef struct SStreamMeta {
Y
yihaoDeng 已提交
368 369 370 371 372 373 374 375 376 377 378 379 380 381 382 383
  char*         path;
  TDB*          db;
  TTB*          pTaskDb;
  TTB*          pCheckpointDb;
  SHashObj*     pTasks;
  SArray*       pTaskList;  // SArray<task_id*>
  void*         ahandle;
  TXN*          txn;
  FTaskExpand*  expandFunc;
  int32_t       vgId;
  SRWLatch      lock;
  int32_t       walScanCounter;
  void*         streamBackend;
  int64_t       streamBackendRid;
  SHashObj*     pTaskBackendUnique;
  TdThreadMutex backendMutex;
384 385
} SStreamMeta;

L
Liu Jicong 已提交
386 387 388
int32_t tEncodeStreamEpInfo(SEncoder* pEncoder, const SStreamChildEpInfo* pInfo);
int32_t tDecodeStreamEpInfo(SDecoder* pDecoder, SStreamChildEpInfo* pInfo);

389 390
SStreamTask* tNewStreamTask(int64_t streamId, int8_t taskLevel, int8_t fillHistory, int64_t triggerParam,
                            SArray* pTaskList);
391 392
int32_t      tEncodeStreamTask(SEncoder* pEncoder, const SStreamTask* pTask);
int32_t      tDecodeStreamTask(SDecoder* pDecoder, SStreamTask* pTask);
393
void         tFreeStreamTask(SStreamTask* pTask);
394
int32_t      tAppendDataToInputQueue(SStreamTask* pTask, SStreamQueueItem* pItem);
395
bool         tInputQueueIsFull(const SStreamTask* pTask);
L
Liu Jicong 已提交
396 397 398 399 400 401 402

typedef struct {
  SMsgHead head;
  int64_t  streamId;
  int32_t  taskId;
} SStreamTaskRunReq;

L
Liu Jicong 已提交
403 404
typedef struct {
  int64_t streamId;
405
  int32_t type;
L
Liu Jicong 已提交
406
  int32_t taskId;
407
  int32_t srcVgId;
L
Liu Jicong 已提交
408 409
  int32_t upstreamTaskId;
  int32_t upstreamChildId;
L
Liu Jicong 已提交
410
  int32_t upstreamNodeId;
L
Liu Jicong 已提交
411
  int32_t blockNum;
412
  int64_t totalLen;
L
Liu Jicong 已提交
413 414
  SArray* dataLen;  // SArray<int32_t>
  SArray* data;     // SArray<SRetrieveTableRsp*>
L
Liu Jicong 已提交
415 416 417 418
} SStreamDispatchReq;

typedef struct {
  int64_t streamId;
419 420 421 422
  int32_t upstreamNodeId;
  int32_t upstreamTaskId;
  int32_t downstreamNodeId;
  int32_t downstreamTaskId;
L
Liu Jicong 已提交
423 424 425
  int8_t  inputStatus;
} SStreamDispatchRsp;

L
Liu Jicong 已提交
426 427
typedef struct {
  int64_t            streamId;
L
Liu Jicong 已提交
428
  int64_t            reqId;
L
Liu Jicong 已提交
429 430 431 432 433 434 435 436 437 438 439 440 441 442 443
  int32_t            srcTaskId;
  int32_t            srcNodeId;
  int32_t            dstTaskId;
  int32_t            dstNodeId;
  int32_t            retrieveLen;
  SRetrieveTableRsp* pRetrieve;
} SStreamRetrieveReq;

typedef struct {
  int64_t streamId;
  int32_t childId;
  int32_t rspFromTaskId;
  int32_t rspToTaskId;
} SStreamRetrieveRsp;

444
typedef struct {
445
  int64_t reqId;
446
  int64_t streamId;
447 448 449 450
  int32_t upstreamNodeId;
  int32_t upstreamTaskId;
  int32_t downstreamNodeId;
  int32_t downstreamTaskId;
451
  int32_t childId;
452
} SStreamTaskCheckReq;
453

L
Liu Jicong 已提交
454
typedef struct {
455
  int64_t reqId;
L
Liu Jicong 已提交
456
  int64_t streamId;
L
Liu Jicong 已提交
457
  int32_t upstreamNodeId;
458 459 460 461 462 463
  int32_t upstreamTaskId;
  int32_t downstreamNodeId;
  int32_t downstreamTaskId;
  int32_t childId;
  int8_t  status;
} SStreamTaskCheckRsp;
L
Liu Jicong 已提交
464 465

typedef struct {
466 467 468
  SMsgHead msgHead;
  int64_t  streamId;
  int32_t  taskId;
L
liuyao 已提交
469 470
  int8_t   igUntreated;
} SStreamScanHistoryReq;
L
Liu Jicong 已提交
471 472 473

typedef struct {
  int64_t streamId;
474 475 476
  int32_t upstreamTaskId;
  int32_t downstreamTaskId;
  int32_t upstreamNodeId;
477
  int32_t childId;
478
} SStreamScanHistoryFinishReq, SStreamTransferReq;
L
Liu Jicong 已提交
479

480 481
int32_t tEncodeStreamScanHistoryFinishReq(SEncoder* pEncoder, const SStreamScanHistoryFinishReq* pReq);
int32_t tDecodeStreamScanHistoryFinishReq(SDecoder* pDecoder, SStreamScanHistoryFinishReq* pReq);
L
Liu Jicong 已提交
482

483 484 485 486 487 488 489 490 491 492 493 494 495 496 497 498 499 500 501 502 503 504 505 506 507 508 509 510 511 512 513 514 515 516 517 518 519 520 521 522 523 524 525 526 527 528 529 530 531 532 533 534 535 536
typedef struct {
  int64_t streamId;
  int64_t checkpointId;
  int32_t taskId;
  int32_t nodeId;
  int64_t expireTime;
} SStreamCheckpointSourceReq;

typedef struct {
  int64_t streamId;
  int64_t checkpointId;
  int32_t taskId;
  int32_t nodeId;
  int64_t expireTime;
} SStreamCheckpointSourceRsp;

int32_t tEncodeSStreamCheckpointSourceReq(SEncoder* pEncoder, const SStreamCheckpointSourceReq* pReq);
int32_t tDecodeSStreamCheckpointSourceReq(SDecoder* pDecoder, SStreamCheckpointSourceReq* pReq);

int32_t tEncodeSStreamCheckpointSourceRsp(SEncoder* pEncoder, const SStreamCheckpointSourceRsp* pRsp);
int32_t tDecodeSStreamCheckpointSourceRsp(SDecoder* pDecoder, SStreamCheckpointSourceRsp* pRsp);

typedef struct {
  SMsgHead msgHead;
  int64_t  streamId;
  int64_t  checkpointId;
  int32_t  downstreamTaskId;
  int32_t  downstreamNodeId;
  int32_t  upstreamTaskId;
  int32_t  upstreamNodeId;
  int32_t  childId;
  int64_t  expireTime;
  int8_t   taskLevel;
} SStreamCheckpointReq;

typedef struct {
  SMsgHead msgHead;
  int64_t  streamId;
  int64_t  checkpointId;
  int32_t  downstreamTaskId;
  int32_t  downstreamNodeId;
  int32_t  upstreamTaskId;
  int32_t  upstreamNodeId;
  int32_t  childId;
  int64_t  expireTime;
  int8_t   taskLevel;
} SStreamCheckpointRsp;

int32_t tEncodeSStreamCheckpointReq(SEncoder* pEncoder, const SStreamCheckpointReq* pReq);
int32_t tDecodeSStreamCheckpointReq(SDecoder* pDecoder, SStreamCheckpointReq* pReq);

int32_t tEncodeSStreamCheckpointRsp(SEncoder* pEncoder, const SStreamCheckpointRsp* pRsp);
int32_t tDecodeSStreamCheckpointRsp(SDecoder* pDecoder, SStreamCheckpointRsp* pRsp);

537 538 539 540 541 542 543 544 545 546 547
typedef struct {
  int64_t streamId;
  int32_t upstreamTaskId;
  int32_t upstreamNodeId;
  int32_t downstreamId;
  int32_t downstreamNode;
} SStreamCompleteHistoryMsg;

int32_t tEncodeCompleteHistoryDataMsg(SEncoder* pEncoder, const SStreamCompleteHistoryMsg* pReq);
int32_t tDecodeCompleteHistoryDataMsg(SDecoder* pDecoder, SStreamCompleteHistoryMsg* pReq);

L
Liu Jicong 已提交
548 549
typedef struct {
  int64_t streamId;
L
Liu Jicong 已提交
550 551 552
  int32_t downstreamTaskId;
  int32_t taskId;
} SStreamRecoverDownstreamReq;
L
Liu Jicong 已提交
553 554

typedef struct {
L
Liu Jicong 已提交
555 556 557 558 559 560
  int64_t streamId;
  int32_t downstreamTaskId;
  int32_t taskId;
  SArray* checkpointVer;  // SArray<SStreamCheckpointInfo>
} SStreamRecoverDownstreamRsp;

561 562
int32_t tEncodeStreamTaskCheckReq(SEncoder* pEncoder, const SStreamTaskCheckReq* pReq);
int32_t tDecodeStreamTaskCheckReq(SDecoder* pDecoder, SStreamTaskCheckReq* pReq);
563

564 565
int32_t tEncodeStreamTaskCheckRsp(SEncoder* pEncoder, const SStreamTaskCheckRsp* pRsp);
int32_t tDecodeStreamTaskCheckRsp(SDecoder* pDecoder, SStreamTaskCheckRsp* pRsp);
566

567 568
int32_t tEncodeSStreamTaskScanHistoryReq(SEncoder* pEncoder, const SStreamRecoverDownstreamReq* pReq);
int32_t tDecodeSStreamTaskScanHistoryReq(SDecoder* pDecoder, SStreamRecoverDownstreamReq* pReq);
L
Liu Jicong 已提交
569 570 571 572

int32_t tEncodeSStreamTaskRecoverRsp(SEncoder* pEncoder, const SStreamRecoverDownstreamRsp* pRsp);
int32_t tDecodeSStreamTaskRecoverRsp(SDecoder* pDecoder, SStreamRecoverDownstreamRsp* pRsp);

573
int32_t tDecodeStreamDispatchReq(SDecoder* pDecoder, SStreamDispatchReq* pReq);
L
Liu Jicong 已提交
574
int32_t tDecodeStreamRetrieveReq(SDecoder* pDecoder, SStreamRetrieveReq* pReq);
L
Liu Jicong 已提交
575 576
void    tDeleteStreamRetrieveReq(SStreamRetrieveReq* pReq);

577
void    tDeleteStreamDispatchReq(SStreamDispatchReq* pReq);
578

579
int32_t streamSetupScheduleTrigger(SStreamTask* pTask);
L
Liu Jicong 已提交
580

L
Liu Jicong 已提交
581
int32_t streamProcessRunReq(SStreamTask* pTask);
582
int32_t streamProcessDispatchMsg(SStreamTask* pTask, SStreamDispatchReq* pReq, SRpcMsg* pMsg, bool exec);
583
int32_t streamProcessDispatchRsp(SStreamTask* pTask, SStreamDispatchRsp* pRsp, int32_t code);
584
void streamTaskCloseUpstreamInput(SStreamTask* pTask, int32_t taskId);
585
void streamTaskOpenAllUpstreamInput(SStreamTask* pTask);
L
Liu Jicong 已提交
586

L
Liu Jicong 已提交
587 588
int32_t streamProcessRetrieveReq(SStreamTask* pTask, SStreamRetrieveReq* pReq, SRpcMsg* pMsg);

589
void    streamTaskInputFail(SStreamTask* pTask);
L
Liu Jicong 已提交
590 591
int32_t streamTryExec(SStreamTask* pTask);
int32_t streamSchedExec(SStreamTask* pTask);
592
int32_t streamTaskOutputResultBlock(SStreamTask* pTask, SStreamDataBlock* pBlock);
593
bool    streamTaskShouldStop(const SStreamStatus* pStatus);
L
liuyao 已提交
594
bool    streamTaskShouldPause(const SStreamStatus* pStatus);
595
bool    streamTaskIsIdle(const SStreamTask* pTask);
L
Liu Jicong 已提交
596

Y
yihaoDeng 已提交
597 598
SStreamChildEpInfo* streamTaskGetUpstreamTaskEpInfo(SStreamTask* pTask, int32_t taskId);
int32_t             streamScanExec(SStreamTask* pTask, int32_t batchSize);
599

Y
yihaoDeng 已提交
600
char* createStreamTaskIdStr(int64_t streamId, int32_t taskId);
601

602
// recover and fill history
603 604
void    streamTaskCheckDownstreamTasks(SStreamTask* pTask);
int32_t streamTaskDoCheckDownstreamTasks(SStreamTask* pTask);
605
int32_t streamTaskLaunchScanHistory(SStreamTask* pTask);
606
int32_t streamTaskCheckStatus(SStreamTask* pTask);
607 608
int32_t streamSendCheckRsp(const SStreamMeta* pMeta, const SStreamTaskCheckReq* pReq, SStreamTaskCheckRsp* pRsp,
                           SRpcHandleInfo* pRpcInfo, int32_t taskId);
609
int32_t streamProcessCheckRsp(SStreamTask* pTask, const SStreamTaskCheckRsp* pRsp);
610
int32_t streamLaunchFillHistoryTask(SStreamTask* pTask);
611
int32_t streamTaskScanHistoryDataComplete(SStreamTask* pTask);
612
int32_t streamStartScanHistoryAsync(SStreamTask* pTask, int8_t igUntreated);
613
bool    streamHistoryTaskSetVerRangeStep2(SStreamTask* pTask, int64_t latestVer);
614

615
// common
Y
yihaoDeng 已提交
616 617
int32_t     streamRestoreParam(SStreamTask* pTask);
int32_t     streamSetStatusNormal(SStreamTask* pTask);
618
const char* streamGetTaskStatusStr(int32_t status);
619
void        streamTaskPause(SStreamTask* pTask);
620 621 622
void        streamTaskResume(SStreamTask* pTask);
void        streamTaskHalt(SStreamTask* pTask);
void        streamTaskResumeFromHalt(SStreamTask* pTask);
623 624
void        streamTaskDisablePause(SStreamTask* pTask);
void        streamTaskEnablePause(SStreamTask* pTask);
625

626
// source level
L
liuyao 已提交
627 628
int32_t streamSetParamForStreamScannerStep1(SStreamTask* pTask, SVersionRange* pVerRange, STimeWindow* pWindow);
int32_t streamSetParamForStreamScannerStep2(SStreamTask* pTask, SVersionRange* pVerRange, STimeWindow* pWindow);
629
int32_t streamSourceScanHistoryData(SStreamTask* pTask);
630
int32_t streamDispatchScanHistoryFinishMsg(SStreamTask* pTask);
H
Haojun Liao 已提交
631

632
int32_t appendTranstateIntoInputQ(SStreamTask* pTask);
H
Haojun Liao 已提交
633

634
// agg level
635
int32_t streamTaskScanHistoryPrepare(SStreamTask* pTask);
Y
yihaoDeng 已提交
636 637
int32_t streamProcessScanHistoryFinishReq(SStreamTask* pTask, SStreamScanHistoryFinishReq* pReq,
                                          SRpcHandleInfo* pRpcInfo);
638
int32_t streamProcessScanHistoryFinishRsp(SStreamTask* pTask);
639

640
// stream task meta
dengyihao's avatar
dengyihao 已提交
641 642
void         streamMetaInit();
void         streamMetaCleanup();
643
SStreamMeta* streamMetaOpen(const char* path, void* ahandle, FTaskExpand expandFunc, int32_t vgId);
L
Liu Jicong 已提交
644
void         streamMetaClose(SStreamMeta* streamMeta);
645 646

// save to b-tree meta store
647
int32_t      streamMetaSaveTask(SStreamMeta* pMeta, SStreamTask* pTask);
648
int32_t      streamMetaRemoveTask(SStreamMeta* pMeta, int32_t taskId);
649
int32_t      streamMetaRegisterTask(SStreamMeta* pMeta, int64_t ver, SStreamTask* pTask, bool* pAdded);
H
Haojun Liao 已提交
650
int32_t      streamMetaUnregisterTask(SStreamMeta* pMeta, int64_t streamId, int32_t taskId);
Y
yihaoDeng 已提交
651
int32_t      streamMetaGetNumOfTasks(SStreamMeta* pMeta);  // todo remove it
652
SStreamTask* streamMetaAcquireTask(SStreamMeta* pMeta, int64_t streamId, int32_t taskId);
L
Liu Jicong 已提交
653 654
void         streamMetaReleaseTask(SStreamMeta* pMeta, SStreamTask* pTask);

655 656 657
int32_t streamMetaBegin(SStreamMeta* pMeta);
int32_t streamMetaCommit(SStreamMeta* pMeta);
int32_t streamLoadTasks(SStreamMeta* pMeta, int64_t ver);
L
Liu Jicong 已提交
658

659 660 661 662 663
// checkpoint
int32_t streamProcessCheckpointSourceReq(SStreamMeta* pMeta, SStreamTask* pTask, SStreamCheckpointSourceReq* pReq);
int32_t streamProcessCheckpointReq(SStreamMeta* pMeta, SStreamTask* pTask, SStreamCheckpointReq* pReq);
int32_t streamProcessCheckpointRsp(SStreamMeta* pMeta, SStreamTask* pTask, SStreamCheckpointRsp* pRsp);

L
liuyao 已提交
664 665
int32_t streamTaskReleaseState(SStreamTask* pTask);
int32_t streamTaskReloadState(SStreamTask* pTask);
H
Haojun Liao 已提交
666 667
int32_t streamAlignTransferState(SStreamTask* pTask);

L
Liu Jicong 已提交
668 669 670 671
#ifdef __cplusplus
}
#endif

L
Liu Jicong 已提交
672
#endif /* ifndef _STREAM_H_ */