tstream.h 17.8 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;

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

L
Liu Jicong 已提交
42
enum {
L
Liu Jicong 已提交
43 44
  TASK_STATUS__NORMAL = 0,
  TASK_STATUS__DROPPING,
L
Liu Jicong 已提交
45 46
  TASK_STATUS__FAIL,
  TASK_STATUS__STOP,
47
  TASK_STATUS__WAIT_DOWNSTREAM,
48 49
  TASK_STATUS__RECOVER_PREPARE,
  TASK_STATUS__RECOVER1,
L
liuyao 已提交
50
  TASK_STATUS__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;
L
Liu Jicong 已提交
125
  int32_t childId;
L
Liu Jicong 已提交
126
  int64_t sourceVer;
L
Liu Jicong 已提交
127
  int64_t reqId;
L
Liu Jicong 已提交
128

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

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

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

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

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

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

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

void streamFreeQitem(SStreamQueueItem* data);

5
54liuyao 已提交
161
#if 0
L
Liu Jicong 已提交
162 163 164 165 166 167
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 已提交
168
#endif
L
Liu Jicong 已提交
169 170 171 172 173

typedef struct {
  SStreamQueueNode* pHead;
} SStreamQueue1;

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

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

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

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

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

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

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

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

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

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

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

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

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

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

typedef struct {
  int8_t reserved;
} STaskSinkFetch;

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

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

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

269
typedef struct SStreamStatus {
270
  int8_t taskStatus;
L
Liu Jicong 已提交
271
  int8_t schedStatus;
L
liuyao 已提交
272
  int8_t keepTaskStatus;
H
Haojun Liao 已提交
273
  bool   transferState;
274
} SStreamStatus;
L
Liu Jicong 已提交
275

276
typedef struct SHistDataRange {
277 278
  SVersionRange range;
  STimeWindow   window;
279
} SHistDataRange;
280

281 282 283 284
typedef struct SSTaskBasicInfo {
  int32_t         nodeId;       // vgroup id or snode id
  SEpSet          epSet;
  int32_t         selfChildId;
285 286
  int32_t         totalLevel;
  int8_t          taskLevel;
287 288
  int8_t          fillHistory;  // is fill history task or not
} SSTaskBasicInfo;
289

290 291 292
typedef struct SDispatchMsgInfo {
  void*   pData;           // current dispatch data
  int16_t msgType;         // dispatch msg type
293 294
  int32_t retryCount;      // retry send data count
  int64_t blockingTs;      // output blocking timestamp
295
} SDispatchMsgInfo;
296

297 298 299 300 301 302
typedef struct {
  int8_t        outputType;
  int8_t        outputStatus;
  SStreamQueue* outputQueue;
} SSTaskOutputInfo;

303 304 305 306 307 308 309 310 311 312
struct SStreamTask {
  SStreamId        id;
  SSTaskBasicInfo  info;
  int8_t           outputType;
  SDispatchMsgInfo msgInfo;
  SStreamStatus    status;
  SCheckpointInfo  chkInfo;
  STaskExec        exec;
  SHistDataRange   dataRange;
  SStreamId        historyTaskId;
H
Haojun Liao 已提交
313
  SStreamId        streamTaskId;
314 315 316
  SArray*          pUpstreamEpInfoList;  // SArray<SStreamChildEpInfo*>, // children info
  int32_t          nextCheckId;
  SArray*          checkpointInfo;  // SArray<SStreamCheckpointInfo>
L
Liu Jicong 已提交
317

318
  // output
L
Liu Jicong 已提交
319 320 321
  union {
    STaskDispatcherFixedEp fixedEpDispatcher;
    STaskDispatcherShuffle shuffleDispatcher;
322 323 324
    STaskSinkTb            tbSink;
    STaskSinkSma           smaSink;
    STaskSinkFetch         fetchSink;
L
Liu Jicong 已提交
325 326
  };

327 328
  int8_t        inputStatus;
  int8_t        outputStatus;
L
Liu Jicong 已提交
329 330
  SStreamQueue* inputQueue;
  SStreamQueue* outputQueue;
L
Liu Jicong 已提交
331

332
  // trigger
333 334 335 336 337
  int8_t        triggerStatus;
  int64_t       triggerParam;
  void*         timer;
  SMsgCb*       pMsgCb;  // msg handle
  SStreamState* pState;  // state backend
338

339 340
  // the followings attributes don't be serialized
  int32_t             recoverTryingDownstream;
341
  int32_t             numOfWaitingUpstream;
342 343 344 345 346 347
  int64_t             checkReqId;
  SArray*             checkReqIds;  // shuffle
  int32_t             refCnt;
  int64_t             checkpointingId;
  int32_t             checkpointAlignCnt;
  struct SStreamMeta* pMeta;
L
liuyao 已提交
348
  SSHashObj*          pNameMap;
349
};
L
Liu Jicong 已提交
350

351 352
// meta
typedef struct SStreamMeta {
353 354 355 356 357
  char*        path;
  TDB*         db;
  TTB*         pTaskDb;
  TTB*         pCheckpointDb;
  SHashObj*    pTasks;
dengyihao's avatar
dengyihao 已提交
358
  SArray*      pTaskList;  // SArray<task_id*>
359 360 361 362 363
  void*        ahandle;
  TXN*         txn;
  FTaskExpand* expandFunc;
  int32_t      vgId;
  SRWLatch     lock;
364
  int32_t      walScanCounter;
dengyihao's avatar
dengyihao 已提交
365
  void*        streamBackend;
dengyihao's avatar
dengyihao 已提交
366
  int64_t      streamBackendRid;
367
  SHashObj*    pTaskBackendUnique;
368 369
} SStreamMeta;

L
Liu Jicong 已提交
370 371 372
int32_t tEncodeStreamEpInfo(SEncoder* pEncoder, const SStreamChildEpInfo* pInfo);
int32_t tDecodeStreamEpInfo(SDecoder* pDecoder, SStreamChildEpInfo* pInfo);

373 374
SStreamTask* tNewStreamTask(int64_t streamId, int8_t taskLevel, int8_t fillHistory, int64_t triggerParam,
                            SArray* pTaskList);
375 376
int32_t      tEncodeStreamTask(SEncoder* pEncoder, const SStreamTask* pTask);
int32_t      tDecodeStreamTask(SDecoder* pDecoder, SStreamTask* pTask);
377
void         tFreeStreamTask(SStreamTask* pTask);
378
int32_t      tAppendDataToInputQueue(SStreamTask* pTask, SStreamQueueItem* pItem);
379
bool         tInputQueueIsFull(const SStreamTask* pTask);
L
Liu Jicong 已提交
380 381 382 383 384 385 386

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

L
Liu Jicong 已提交
387 388 389
typedef struct {
  int64_t streamId;
  int32_t taskId;
L
Liu Jicong 已提交
390 391 392
  int32_t dataSrcVgId;
  int32_t upstreamTaskId;
  int32_t upstreamChildId;
L
Liu Jicong 已提交
393
  int32_t upstreamNodeId;
L
Liu Jicong 已提交
394
  int32_t blockNum;
395
  int64_t totalLen;
L
Liu Jicong 已提交
396 397
  SArray* dataLen;  // SArray<int32_t>
  SArray* data;     // SArray<SRetrieveTableRsp*>
L
Liu Jicong 已提交
398 399 400 401
} SStreamDispatchReq;

typedef struct {
  int64_t streamId;
402 403 404 405
  int32_t upstreamNodeId;
  int32_t upstreamTaskId;
  int32_t downstreamNodeId;
  int32_t downstreamTaskId;
L
Liu Jicong 已提交
406 407 408
  int8_t  inputStatus;
} SStreamDispatchRsp;

L
Liu Jicong 已提交
409 410
typedef struct {
  int64_t            streamId;
L
Liu Jicong 已提交
411
  int64_t            reqId;
L
Liu Jicong 已提交
412 413 414 415 416 417 418 419 420 421 422 423 424 425 426
  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;

427
typedef struct {
428
  int64_t reqId;
429
  int64_t streamId;
430 431 432 433
  int32_t upstreamNodeId;
  int32_t upstreamTaskId;
  int32_t downstreamNodeId;
  int32_t downstreamTaskId;
434
  int32_t childId;
435
} SStreamTaskCheckReq;
436

L
Liu Jicong 已提交
437
typedef struct {
438
  int64_t reqId;
L
Liu Jicong 已提交
439
  int64_t streamId;
L
Liu Jicong 已提交
440
  int32_t upstreamNodeId;
441 442 443 444 445 446
  int32_t upstreamTaskId;
  int32_t downstreamNodeId;
  int32_t downstreamTaskId;
  int32_t childId;
  int8_t  status;
} SStreamTaskCheckRsp;
L
Liu Jicong 已提交
447 448

typedef struct {
449 450 451 452
  SMsgHead msgHead;
  int64_t  streamId;
  int32_t  taskId;
} SStreamRecoverStep1Req, SStreamRecoverStep2Req;
L
Liu Jicong 已提交
453 454 455 456

typedef struct {
  int64_t streamId;
  int32_t taskId;
457
  int32_t childId;
H
Haojun Liao 已提交
458
} SStreamRecoverFinishReq, SStreamTransferReq;
L
Liu Jicong 已提交
459

H
Haojun Liao 已提交
460 461
int32_t tEncodeStreamRecoverFinishReq(SEncoder* pEncoder, const SStreamRecoverFinishReq* pReq);
int32_t tDecodeStreamRecoverFinishReq(SDecoder* pDecoder, SStreamRecoverFinishReq* pReq);
L
Liu Jicong 已提交
462

463 464 465 466 467 468 469 470 471 472 473 474 475 476 477 478 479 480 481 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
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);

L
Liu Jicong 已提交
517 518
typedef struct {
  int64_t streamId;
L
Liu Jicong 已提交
519 520 521
  int32_t downstreamTaskId;
  int32_t taskId;
} SStreamRecoverDownstreamReq;
L
Liu Jicong 已提交
522 523

typedef struct {
L
Liu Jicong 已提交
524 525 526 527 528 529
  int64_t streamId;
  int32_t downstreamTaskId;
  int32_t taskId;
  SArray* checkpointVer;  // SArray<SStreamCheckpointInfo>
} SStreamRecoverDownstreamRsp;

530 531 532 533 534 535
int32_t tEncodeSStreamTaskCheckReq(SEncoder* pEncoder, const SStreamTaskCheckReq* pReq);
int32_t tDecodeSStreamTaskCheckReq(SDecoder* pDecoder, SStreamTaskCheckReq* pReq);

int32_t tEncodeSStreamTaskCheckRsp(SEncoder* pEncoder, const SStreamTaskCheckRsp* pRsp);
int32_t tDecodeSStreamTaskCheckRsp(SDecoder* pDecoder, SStreamTaskCheckRsp* pRsp);

L
Liu Jicong 已提交
536
int32_t tEncodeSStreamTaskRecoverReq(SEncoder* pEncoder, const SStreamRecoverDownstreamReq* pReq);
L
Liu Jicong 已提交
537 538 539 540 541
int32_t tDecodeSStreamTaskRecoverReq(SDecoder* pDecoder, SStreamRecoverDownstreamReq* pReq);

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

542
int32_t tDecodeStreamDispatchReq(SDecoder* pDecoder, SStreamDispatchReq* pReq);
L
Liu Jicong 已提交
543
int32_t tDecodeStreamRetrieveReq(SDecoder* pDecoder, SStreamRetrieveReq* pReq);
L
Liu Jicong 已提交
544 545
void    tDeleteStreamRetrieveReq(SStreamRetrieveReq* pReq);

546 547 548
int32_t tInitStreamDispatchReq(SStreamDispatchReq* pReq, const SStreamTask* pTask, int32_t vgId, int32_t numOfBlocks,
                               int64_t dstTaskId);
void    tDeleteStreamDispatchReq(SStreamDispatchReq* pReq);
549

550
int32_t streamSetupTrigger(SStreamTask* pTask);
L
Liu Jicong 已提交
551

L
Liu Jicong 已提交
552
int32_t streamProcessRunReq(SStreamTask* pTask);
553
int32_t streamProcessDispatchMsg(SStreamTask* pTask, SStreamDispatchReq* pReq, SRpcMsg* pMsg, bool exec);
554
int32_t streamProcessDispatchRsp(SStreamTask* pTask, SStreamDispatchRsp* pRsp, int32_t code);
L
Liu Jicong 已提交
555

L
Liu Jicong 已提交
556
int32_t streamProcessRetrieveReq(SStreamTask* pTask, SStreamRetrieveReq* pReq, SRpcMsg* pMsg);
5
54liuyao 已提交
557
// int32_t streamProcessRetrieveRsp(SStreamTask* pTask, SStreamRetrieveRsp* pRsp);
L
Liu Jicong 已提交
558

559
void    streamTaskInputFail(SStreamTask* pTask);
L
Liu Jicong 已提交
560 561
int32_t streamTryExec(SStreamTask* pTask);
int32_t streamSchedExec(SStreamTask* pTask);
562
int32_t streamTaskOutputResultBlock(SStreamTask* pTask, SStreamDataBlock* pBlock);
563
bool    streamTaskShouldStop(const SStreamStatus* pStatus);
L
liuyao 已提交
564
bool    streamTaskShouldPause(const SStreamStatus* pStatus);
L
Liu Jicong 已提交
565

566 567 568
int32_t streamScanExec(SStreamTask* pTask, int32_t batchSz);

// recover and fill history
569
int32_t streamTaskCheckDownstreamTasks(SStreamTask* pTask);
570
int32_t streamTaskLaunchRecover(SStreamTask* pTask);
571
int32_t streamTaskCheckStatus(SStreamTask* pTask);
572 573
int32_t streamProcessCheckRsp(SStreamTask* pTask, const SStreamTaskCheckRsp* pRsp);
int32_t streamTaskStartHistoryTask(SStreamTask* pTask, int64_t ver);
574
int32_t streamTaskScanHistoryDataComplete(SStreamTask* pTask);
575

576 577 578 579 580
// common
int32_t streamSetParamForRecover(SStreamTask* pTask);
int32_t streamRestoreParam(SStreamTask* pTask);
int32_t streamSetStatusNormal(SStreamTask* pTask);
// source level
581
int32_t streamSourceRecoverPrepareStep1(SStreamTask* pTask, SVersionRange *pVerRange, STimeWindow* pWindow);
582 583 584 585
int32_t streamBuildSourceRecover1Req(SStreamTask* pTask, SStreamRecoverStep1Req* pReq);
int32_t streamSourceRecoverScanStep1(SStreamTask* pTask);
int32_t streamBuildSourceRecover2Req(SStreamTask* pTask, SStreamRecoverStep2Req* pReq);
int32_t streamSourceRecoverScanStep2(SStreamTask* pTask, int64_t ver);
586
int32_t streamDispatchRecoverFinishMsg(SStreamTask* pTask);
H
Haojun Liao 已提交
587 588 589

int32_t streamDispatchTransferStateMsg(SStreamTask* pTask);

590 591 592 593
// agg level
int32_t streamAggRecoverPrepare(SStreamTask* pTask);
int32_t streamProcessRecoverFinishReq(SStreamTask* pTask, int32_t childId);

dengyihao's avatar
dengyihao 已提交
594 595
void         streamMetaInit();
void         streamMetaCleanup();
596
SStreamMeta* streamMetaOpen(const char* path, void* ahandle, FTaskExpand expandFunc, int32_t vgId);
L
Liu Jicong 已提交
597 598
void         streamMetaClose(SStreamMeta* streamMeta);

599 600 601 602
int32_t streamMetaSaveTask(SStreamMeta* pMeta, SStreamTask* pTask);
int32_t streamMetaAddDeployedTask(SStreamMeta* pMeta, int64_t ver, SStreamTask* pTask);
int32_t streamMetaAddSerializedTask(SStreamMeta* pMeta, int64_t checkpointVer, char* msg, int32_t msgLen);
int32_t streamMetaGetNumOfTasks(const SStreamMeta* pMeta);
L
Liu Jicong 已提交
603

L
Liu Jicong 已提交
604 605
SStreamTask* streamMetaAcquireTask(SStreamMeta* pMeta, int32_t taskId);
void         streamMetaReleaseTask(SStreamMeta* pMeta, SStreamTask* pTask);
606
void         streamMetaRemoveTask(SStreamMeta* pMeta, int32_t taskId);
L
Liu Jicong 已提交
607

608 609 610 611
int32_t streamMetaBegin(SStreamMeta* pMeta);
int32_t streamMetaCommit(SStreamMeta* pMeta);
int32_t streamMetaRollBack(SStreamMeta* pMeta);
int32_t streamLoadTasks(SStreamMeta* pMeta, int64_t ver);
L
Liu Jicong 已提交
612

613 614 615 616 617
// 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
Liu Jicong 已提交
618 619 620 621
#ifdef __cplusplus
}
#endif

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