tstream.h 19.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,
Y
yihaoDeng 已提交
47
  TASK_STATUS__SCAN_HISTORY,  // stream task scan history data by using tsdbread in the stream scanner
48
  TASK_STATUS__HALT,          // stream task will handle all data in the input queue, and then paused, todo remove it?
49
  TASK_STATUS__PAUSE,         // pause
L
Liu Jicong 已提交
50 51 52
};

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

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

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

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

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

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

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

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

103 104
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 已提交
105 106

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

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

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

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

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

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

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

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

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

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

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

void streamFreeQitem(SStreamQueueItem* data);

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

typedef struct {
  SStreamQueueNode* pHead;
} SStreamQueue1;

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

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

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

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

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

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

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

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

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

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

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

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

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

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

typedef struct {
  int8_t reserved;
} STaskSinkFetch;

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

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

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

268
typedef struct SStreamStatus {
Y
yihaoDeng 已提交
269
  int8_t        taskStatus;
270
  int8_t        downstreamReady; // downstream tasks are all ready now, if this flag is set
Y
yihaoDeng 已提交
271 272 273
  int8_t        schedStatus;
  int8_t        keepTaskStatus;
  bool          transferState;
H
Haojun Liao 已提交
274
  int8_t        timerActive;     // timer is active
275
  int8_t        pauseAllowed;    // allowed task status to be set to be paused
276
} SStreamStatus;
L
Liu Jicong 已提交
277

278
typedef struct SHistDataRange {
279 280
  SVersionRange range;
  STimeWindow   window;
281
} SHistDataRange;
282

283
typedef struct SSTaskBasicInfo {
Y
yihaoDeng 已提交
284 285 286 287 288 289
  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
290
} SSTaskBasicInfo;
291

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

299
typedef struct {
300 301 302 303
  int8_t        type;
  int8_t        status;
  SStreamQueue* queue;
} STaskOutputInfo;
304

305 306 307
struct SStreamTask {
  SStreamId        id;
  SSTaskBasicInfo  info;
308
  STaskOutputInfo  outputInfo;
309 310 311 312 313 314
  SDispatchMsgInfo msgInfo;
  SStreamStatus    status;
  SCheckpointInfo  chkInfo;
  STaskExec        exec;
  SHistDataRange   dataRange;
  SStreamId        historyTaskId;
H
Haojun Liao 已提交
315
  SStreamId        streamTaskId;
316 317 318
  SArray*          pUpstreamEpInfoList;  // SArray<SStreamChildEpInfo*>, // children info
  int32_t          nextCheckId;
  SArray*          checkpointInfo;  // SArray<SStreamCheckpointInfo>
319
  int64_t          initTs;
320
  // output
L
Liu Jicong 已提交
321 322 323
  union {
    STaskDispatcherFixedEp fixedEpDispatcher;
    STaskDispatcherShuffle shuffleDispatcher;
324 325 326
    STaskSinkTb            tbSink;
    STaskSinkSma           smaSink;
    STaskSinkFetch         fetchSink;
L
Liu Jicong 已提交
327 328
  };

329
  int8_t        inputStatus;
L
Liu Jicong 已提交
330
  SStreamQueue* inputQueue;
L
Liu Jicong 已提交
331

332
  // trigger
333 334
  int8_t        triggerStatus;
  int64_t       triggerParam;
335 336
  void*         schedTimer;
  void*         launchTaskTimer;
337 338
  SMsgCb*       pMsgCb;  // msg handle
  SStreamState* pState;  // state backend
339
  SArray*       pRspMsgList;
H
Haojun Liao 已提交
340
  TdThreadMutex lock;
341

342
  // the followings attributes don't be serialized
343
  int32_t             notReadyTasks;
344
  int32_t             numOfWaitingUpstream;
345 346 347 348 349
  int64_t             checkReqId;
  SArray*             checkReqIds;  // shuffle
  int32_t             refCnt;
  int64_t             checkpointingId;
  int32_t             checkpointAlignCnt;
H
Haojun Liao 已提交
350
  int32_t             transferStateAlignCnt;
351
  struct SStreamMeta* pMeta;
L
liuyao 已提交
352
  SSHashObj*          pNameMap;
353
};
L
Liu Jicong 已提交
354

355 356
// meta
typedef struct SStreamMeta {
Y
yihaoDeng 已提交
357 358 359 360 361 362 363 364 365 366 367 368 369 370 371 372
  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;
373 374
} SStreamMeta;

L
Liu Jicong 已提交
375 376 377
int32_t tEncodeStreamEpInfo(SEncoder* pEncoder, const SStreamChildEpInfo* pInfo);
int32_t tDecodeStreamEpInfo(SDecoder* pDecoder, SStreamChildEpInfo* pInfo);

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

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

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

typedef struct {
  int64_t streamId;
407 408 409 410
  int32_t upstreamNodeId;
  int32_t upstreamTaskId;
  int32_t downstreamNodeId;
  int32_t downstreamTaskId;
L
Liu Jicong 已提交
411 412 413
  int8_t  inputStatus;
} SStreamDispatchRsp;

L
Liu Jicong 已提交
414 415
typedef struct {
  int64_t            streamId;
L
Liu Jicong 已提交
416
  int64_t            reqId;
L
Liu Jicong 已提交
417 418 419 420 421 422 423 424 425 426 427 428 429 430 431
  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;

432
typedef struct {
433
  int64_t reqId;
434
  int64_t streamId;
435 436 437 438
  int32_t upstreamNodeId;
  int32_t upstreamTaskId;
  int32_t downstreamNodeId;
  int32_t downstreamTaskId;
439
  int32_t childId;
440
} SStreamTaskCheckReq;
441

L
Liu Jicong 已提交
442
typedef struct {
443
  int64_t reqId;
L
Liu Jicong 已提交
444
  int64_t streamId;
L
Liu Jicong 已提交
445
  int32_t upstreamNodeId;
446 447 448 449 450 451
  int32_t upstreamTaskId;
  int32_t downstreamNodeId;
  int32_t downstreamTaskId;
  int32_t childId;
  int8_t  status;
} SStreamTaskCheckRsp;
L
Liu Jicong 已提交
452 453

typedef struct {
454 455 456
  SMsgHead msgHead;
  int64_t  streamId;
  int32_t  taskId;
L
liuyao 已提交
457 458
  int8_t   igUntreated;
} SStreamScanHistoryReq;
L
Liu Jicong 已提交
459 460 461

typedef struct {
  int64_t streamId;
462 463 464
  int32_t upstreamTaskId;
  int32_t downstreamTaskId;
  int32_t upstreamNodeId;
465
  int32_t childId;
466
} SStreamScanHistoryFinishReq, SStreamTransferReq;
L
Liu Jicong 已提交
467

468 469
int32_t tEncodeStreamScanHistoryFinishReq(SEncoder* pEncoder, const SStreamScanHistoryFinishReq* pReq);
int32_t tDecodeStreamScanHistoryFinishReq(SDecoder* pDecoder, SStreamScanHistoryFinishReq* pReq);
L
Liu Jicong 已提交
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 517 518 519 520 521 522 523 524
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);

525 526 527 528 529 530 531 532 533 534 535
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 已提交
536 537
typedef struct {
  int64_t streamId;
L
Liu Jicong 已提交
538 539 540
  int32_t downstreamTaskId;
  int32_t taskId;
} SStreamRecoverDownstreamReq;
L
Liu Jicong 已提交
541 542

typedef struct {
L
Liu Jicong 已提交
543 544 545 546 547 548
  int64_t streamId;
  int32_t downstreamTaskId;
  int32_t taskId;
  SArray* checkpointVer;  // SArray<SStreamCheckpointInfo>
} SStreamRecoverDownstreamRsp;

549 550
int32_t tEncodeStreamTaskCheckReq(SEncoder* pEncoder, const SStreamTaskCheckReq* pReq);
int32_t tDecodeStreamTaskCheckReq(SDecoder* pDecoder, SStreamTaskCheckReq* pReq);
551

552 553
int32_t tEncodeStreamTaskCheckRsp(SEncoder* pEncoder, const SStreamTaskCheckRsp* pRsp);
int32_t tDecodeStreamTaskCheckRsp(SDecoder* pDecoder, SStreamTaskCheckRsp* pRsp);
554

555 556
int32_t tEncodeSStreamTaskScanHistoryReq(SEncoder* pEncoder, const SStreamRecoverDownstreamReq* pReq);
int32_t tDecodeSStreamTaskScanHistoryReq(SDecoder* pDecoder, SStreamRecoverDownstreamReq* pReq);
L
Liu Jicong 已提交
557 558 559 560

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

561
int32_t tDecodeStreamDispatchReq(SDecoder* pDecoder, SStreamDispatchReq* pReq);
L
Liu Jicong 已提交
562
int32_t tDecodeStreamRetrieveReq(SDecoder* pDecoder, SStreamRetrieveReq* pReq);
L
Liu Jicong 已提交
563 564
void    tDeleteStreamRetrieveReq(SStreamRetrieveReq* pReq);

565 566 567
int32_t tInitStreamDispatchReq(SStreamDispatchReq* pReq, const SStreamTask* pTask, int32_t vgId, int32_t numOfBlocks,
                               int64_t dstTaskId);
void    tDeleteStreamDispatchReq(SStreamDispatchReq* pReq);
568

569
int32_t streamSetupScheduleTrigger(SStreamTask* pTask);
L
Liu Jicong 已提交
570

L
Liu Jicong 已提交
571
int32_t streamProcessRunReq(SStreamTask* pTask);
572
int32_t streamProcessDispatchMsg(SStreamTask* pTask, SStreamDispatchReq* pReq, SRpcMsg* pMsg, bool exec);
573
int32_t streamProcessDispatchRsp(SStreamTask* pTask, SStreamDispatchRsp* pRsp, int32_t code);
L
Liu Jicong 已提交
574

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

577
void    streamTaskInputFail(SStreamTask* pTask);
L
Liu Jicong 已提交
578 579
int32_t streamTryExec(SStreamTask* pTask);
int32_t streamSchedExec(SStreamTask* pTask);
580
int32_t streamTaskOutputResultBlock(SStreamTask* pTask, SStreamDataBlock* pBlock);
581
bool    streamTaskShouldStop(const SStreamStatus* pStatus);
L
liuyao 已提交
582
bool    streamTaskShouldPause(const SStreamStatus* pStatus);
583
bool    streamTaskIsIdle(const SStreamTask* pTask);
L
Liu Jicong 已提交
584

585
SStreamChildEpInfo * streamTaskGetUpstreamTaskEpInfo(SStreamTask* pTask, int32_t taskId);
586 587
int32_t streamScanExec(SStreamTask* pTask, int32_t batchSz);

588 589
char*   createStreamTaskIdStr(int64_t streamId, int32_t taskId);

590
// recover and fill history
591 592
void    streamTaskCheckDownstreamTasks(SStreamTask* pTask);
int32_t streamTaskDoCheckDownstreamTasks(SStreamTask* pTask);
593
int32_t streamTaskLaunchScanHistory(SStreamTask* pTask);
594
int32_t streamTaskCheckStatus(SStreamTask* pTask);
595 596
int32_t streamSendCheckRsp(const SStreamMeta* pMeta, const SStreamTaskCheckReq* pReq, SStreamTaskCheckRsp* pRsp,
                           SRpcHandleInfo* pRpcInfo, int32_t taskId);
597
int32_t streamProcessCheckRsp(SStreamTask* pTask, const SStreamTaskCheckRsp* pRsp);
598
int32_t streamLaunchFillHistoryTask(SStreamTask* pTask);
599
int32_t streamTaskScanHistoryDataComplete(SStreamTask* pTask);
L
liuyao 已提交
600
int32_t streamStartRecoverTask(SStreamTask* pTask, int8_t igUntreated);
601 602
void    streamHistoryTaskSetVerRangeStep2(SStreamTask* pTask);

L
liuyao 已提交
603 604 605
bool    streamTaskRecoverScanStep1Finished(SStreamTask* pTask);
bool    streamTaskRecoverScanStep2Finished(SStreamTask* pTask);
int32_t streamTaskRecoverSetAllStepFinished(SStreamTask* pTask);
606

607
// common
608
int32_t     streamSetParamForScanHistory(SStreamTask* pTask);
Y
yihaoDeng 已提交
609 610
int32_t     streamRestoreParam(SStreamTask* pTask);
int32_t     streamSetStatusNormal(SStreamTask* pTask);
611
const char* streamGetTaskStatusStr(int32_t status);
612 613 614
void        streamTaskPause(SStreamTask* pTask);
void        streamTaskDisablePause(SStreamTask* pTask);
void        streamTaskEnablePause(SStreamTask* pTask);
615

616
// source level
L
liuyao 已提交
617 618
int32_t streamSetParamForStreamScannerStep1(SStreamTask* pTask, SVersionRange* pVerRange, STimeWindow* pWindow);
int32_t streamSetParamForStreamScannerStep2(SStreamTask* pTask, SVersionRange* pVerRange, STimeWindow* pWindow);
L
liuyao 已提交
619
int32_t streamBuildSourceRecover1Req(SStreamTask* pTask, SStreamScanHistoryReq* pReq, int8_t igUntreated);
620
int32_t streamSourceScanHistoryData(SStreamTask* pTask);
621
int32_t streamDispatchScanHistoryFinishMsg(SStreamTask* pTask);
H
Haojun Liao 已提交
622 623 624

int32_t streamDispatchTransferStateMsg(SStreamTask* pTask);

625
// agg level
626 627 628
int32_t streamTaskScanHistoryPrepare(SStreamTask* pTask);
int32_t streamProcessScanHistoryFinishReq(SStreamTask* pTask, SStreamScanHistoryFinishReq *pReq, SRpcHandleInfo* pRpcInfo);
int32_t streamProcessScanHistoryFinishRsp(SStreamTask* pTask);
629

630
// stream task meta
dengyihao's avatar
dengyihao 已提交
631 632
void         streamMetaInit();
void         streamMetaCleanup();
633
SStreamMeta* streamMetaOpen(const char* path, void* ahandle, FTaskExpand expandFunc, int32_t vgId);
L
Liu Jicong 已提交
634
void         streamMetaClose(SStreamMeta* streamMeta);
635 636 637 638
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);   // todo remove it
L
Liu Jicong 已提交
639 640
SStreamTask* streamMetaAcquireTask(SStreamMeta* pMeta, int32_t taskId);
void         streamMetaReleaseTask(SStreamMeta* pMeta, SStreamTask* pTask);
641
void         streamMetaRemoveTask(SStreamMeta* pMeta, int32_t taskId);
L
Liu Jicong 已提交
642

643 644 645
int32_t streamMetaBegin(SStreamMeta* pMeta);
int32_t streamMetaCommit(SStreamMeta* pMeta);
int32_t streamLoadTasks(SStreamMeta* pMeta, int64_t ver);
L
Liu Jicong 已提交
646

647 648 649 650 651
// 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 已提交
652 653
int32_t streamTaskReleaseState(SStreamTask* pTask);
int32_t streamTaskReloadState(SStreamTask* pTask);
H
Haojun Liao 已提交
654 655
int32_t streamAlignTransferState(SStreamTask* pTask);

L
liuyao 已提交
656

L
Liu Jicong 已提交
657 658 659 660
#ifdef __cplusplus
}
#endif

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