tq.h 5.6 KB
Newer Older
L
Liu Jicong 已提交
1 2 3 4 5 6 7 8 9 10 11 12 13 14
/*
 * 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/>.
 */
H
Hongze Cheng 已提交
15 16 17 18

#ifndef _TD_VNODE_TQ_H_
#define _TD_VNODE_TQ_H_

H
Hongze Cheng 已提交
19 20
#include "vnodeInt.h"

L
Liu Jicong 已提交
21 22 23 24
#include "executor.h"
#include "os.h"
#include "thash.h"
#include "tmsg.h"
25
#include "tqueue.h"
L
Liu Jicong 已提交
26
#include "trpc.h"
L
Liu Jicong 已提交
27 28 29
#include "ttimer.h"
#include "wal.h"

H
Hongze Cheng 已提交
30 31 32 33
#ifdef __cplusplus
extern "C" {
#endif

H
Hongze Cheng 已提交
34 35
// tqDebug ===================
// clang-format off
L
Liu Jicong 已提交
36 37 38 39 40 41
#define tqFatal(...) do { if (tqDebugFlag & DEBUG_FATAL) { taosPrintLog("TQ  FATAL ", DEBUG_FATAL, 255, __VA_ARGS__); }}     while(0)
#define tqError(...) do { if (tqDebugFlag & DEBUG_ERROR) { taosPrintLog("TQ  ERROR ", DEBUG_ERROR, 255, __VA_ARGS__); }}     while(0)
#define tqWarn(...)  do { if (tqDebugFlag & DEBUG_WARN)  { taosPrintLog("TQ  WARN ", DEBUG_WARN, 255, __VA_ARGS__); }}       while(0)
#define tqInfo(...)  do { if (tqDebugFlag & DEBUG_INFO)  { taosPrintLog("TQ  ", DEBUG_INFO, 255, __VA_ARGS__); }}            while(0)
#define tqDebug(...) do { if (tqDebugFlag & DEBUG_DEBUG) { taosPrintLog("TQ  ", DEBUG_DEBUG, tqDebugFlag, __VA_ARGS__); }} while(0)
#define tqTrace(...) do { if (tqDebugFlag & DEBUG_TRACE) { taosPrintLog("TQ  ", DEBUG_TRACE, tqDebugFlag, __VA_ARGS__); }} while(0)
L
Liu Jicong 已提交
42

H
Hongze Cheng 已提交
43 44
// clang-format on

H
Hongze Cheng 已提交
45 46
typedef struct STqOffsetStore STqOffsetStore;

L
Liu Jicong 已提交
47 48
// tqPush

L
Liu Jicong 已提交
49
typedef struct {
L
Liu Jicong 已提交
50
  // msg info
L
Liu Jicong 已提交
51
  int64_t consumerId;
L
Liu Jicong 已提交
52 53
  int64_t reqOffset;
  int64_t processedVer;
L
Liu Jicong 已提交
54
  int32_t epoch;
L
Liu Jicong 已提交
55 56 57
  // rpc info
  int64_t        reqId;
  SRpcHandleInfo rpcInfo;
L
Liu Jicong 已提交
58 59
  tmr_h          timerId;
  int8_t         tmrStopped;
L
Liu Jicong 已提交
60
  // exec
L
Liu Jicong 已提交
61 62 63 64
  int8_t       inputStatus;
  int8_t       execStatus;
  SStreamQueue inputQ;
  SRWLatch     lock;
L
Liu Jicong 已提交
65
} STqPushHandle;
L
Liu Jicong 已提交
66 67 68 69 70 71 72 73

// tqExec

typedef struct {
  char*       qmsg;
} STqExecCol;

typedef struct {
74
  int64_t     suid;
L
Liu Jicong 已提交
75 76 77
} STqExecTb;

typedef struct {
78
  SHashObj*   pFilterOutTbUid;
L
Liu Jicong 已提交
79 80 81 82 83
} STqExecDb;

typedef struct {
  int8_t subType;

L
Liu Jicong 已提交
84
  STqReader* pExecReader;
85
  qTaskInfo_t task;
L
Liu Jicong 已提交
86 87 88 89
  union {
    STqExecCol execCol;
    STqExecTb  execTb;
    STqExecDb  execDb;
L
Liu Jicong 已提交
90
  };
wmmhello's avatar
wmmhello 已提交
91
  int32_t         numOfCols;       // number of out pout column, temporarily used
L
Liu Jicong 已提交
92
  SSchemaWrapper* pSchemaWrapper;  // columns that are involved in query
L
Liu Jicong 已提交
93
} STqExecHandle;
H
Hongze Cheng 已提交
94

L
Liu Jicong 已提交
95 96 97 98 99
typedef struct {
  // info
  char    subKey[TSDB_SUBSCRIBE_KEY_LEN];
  int64_t consumerId;
  int32_t epoch;
L
Liu Jicong 已提交
100
  int8_t  fetchMeta;
L
Liu Jicong 已提交
101

L
Liu Jicong 已提交
102 103
  int64_t snapshotVer;

wmmhello's avatar
wmmhello 已提交
104
  SWalReader* pWalReader;
L
Liu Jicong 已提交
105

106 107
  SWalRef* pRef;

L
Liu Jicong 已提交
108 109 110 111 112
  // push
  STqPushHandle pushHandle;

  // exec
  STqExecHandle execHandle;
L
Liu Jicong 已提交
113

L
Liu Jicong 已提交
114
} STqHandle;
115

H
Hongze Cheng 已提交
116
struct STQ {
L
Liu Jicong 已提交
117 118
  SVnode*   pVnode;
  char*     path;
119 120 121
  SHashObj* pPushMgr;    // consumerId -> STqHandle*
  SHashObj* pHandle;     // subKey -> STqHandle
  SHashObj* pCheckInfo;  // topic -> SAlterCheckInfo
L
Liu Jicong 已提交
122

123
  STqOffsetStore* pOffsetStore;
L
Liu Jicong 已提交
124

125
  TDB* pMetaDB;
L
Liu Jicong 已提交
126
  TTB* pExecStore;
127
  TTB* pCheckStore;
L
Liu Jicong 已提交
128

L
Liu Jicong 已提交
129
  SStreamMeta* pStreamMeta;
H
Hongze Cheng 已提交
130 131 132 133 134 135 136
};

typedef struct {
  int8_t inited;
  tmr_h  timer;
} STqMgmt;

L
Liu Jicong 已提交
137
static STqMgmt tqMgmt = {0};
H
Hongze Cheng 已提交
138

L
Liu Jicong 已提交
139 140 141
int32_t tEncodeSTqHandle(SEncoder* pEncoder, const STqHandle* pHandle);
int32_t tDecodeSTqHandle(SDecoder* pDecoder, STqHandle* pHandle);

L
Liu Jicong 已提交
142
// tqRead
wmmhello's avatar
wmmhello 已提交
143
int64_t tqScan(STQ* pTq, const STqHandle* pHandle, SMqDataRsp* pRsp, SMqMetaRsp* pMetaRsp, STqOffsetVal* offset);
wmmhello's avatar
wmmhello 已提交
144
int64_t tqFetchLog(STQ* pTq, STqHandle* pHandle, int64_t* fetchOffset, SWalCkHead** pHeadWithCkSum);
L
Liu Jicong 已提交
145

wmmhello's avatar
wmmhello 已提交
146 147
// tqExec
int32_t tqLogScanExec(STQ* pTq, STqExecHandle* pExec, SSubmitReq* pReq, SMqDataRsp* pRsp);
L
Liu Jicong 已提交
148
int32_t tqSendDataRsp(STQ* pTq, const SRpcMsg* pMsg, const SMqPollReq* pReq, const SMqDataRsp* pRsp);
L
Liu Jicong 已提交
149

L
Liu Jicong 已提交
150 151 152 153 154
// tqMeta
int32_t tqMetaOpen(STQ* pTq);
int32_t tqMetaClose(STQ* pTq);
int32_t tqMetaSaveHandle(STQ* pTq, const char* key, const STqHandle* pHandle);
int32_t tqMetaDeleteHandle(STQ* pTq, const char* key);
L
Liu Jicong 已提交
155
int32_t tqMetaRestoreHandle(STQ* pTq);
156 157 158
int32_t tqMetaSaveCheckInfo(STQ* pTq, const char* key, const void* value, int32_t vLen);
int32_t tqMetaDeleteCheckInfo(STQ* pTq, const char* key);
int32_t tqMetaRestoreCheckInfo(STQ* pTq);
L
Liu Jicong 已提交
159

160 161 162
typedef struct {
  int32_t size;
} STqOffsetHead;
L
Liu Jicong 已提交
163

164
STqOffsetStore* tqOffsetOpen();
L
Liu Jicong 已提交
165
void            tqOffsetClose(STqOffsetStore*);
166 167
STqOffset*      tqOffsetRead(STqOffsetStore* pStore, const char* subscribeKey);
int32_t         tqOffsetWrite(STqOffsetStore* pStore, const STqOffset* pOffset);
L
Liu Jicong 已提交
168
int32_t         tqOffsetDelete(STqOffsetStore* pStore, const char* subscribeKey);
L
Liu Jicong 已提交
169
int32_t         tqOffsetCommitFile(STqOffsetStore* pStore);
170 171 172

// tqSink
void tqTableSink(SStreamTask* pTask, void* vnode, int64_t ver, void* data);
H
Hongze Cheng 已提交
173

L
Liu Jicong 已提交
174 175 176 177
// tqOffset
char*   tqOffsetBuildFName(const char* path, int32_t ver);
int32_t tqOffsetRestoreFromFile(STqOffsetStore* pStore, const char* fname);

L
Liu Jicong 已提交
178 179 180 181 182 183
static FORCE_INLINE void tqOffsetResetToData(STqOffsetVal* pOffsetVal, int64_t uid, int64_t ts) {
  pOffsetVal->type = TMQ_OFFSET__SNAPSHOT_DATA;
  pOffsetVal->uid = uid;
  pOffsetVal->ts = ts;
}

wmmhello's avatar
wmmhello 已提交
184
static FORCE_INLINE void tqOffsetResetToMeta(STqOffsetVal* pOffsetVal, int64_t uid) {
185
  pOffsetVal->type = TMQ_OFFSET__SNAPSHOT_META;
wmmhello's avatar
wmmhello 已提交
186
  pOffsetVal->uid = uid;
187 188
}

L
Liu Jicong 已提交
189 190 191 192 193
static FORCE_INLINE void tqOffsetResetToLog(STqOffsetVal* pOffsetVal, int64_t ver) {
  pOffsetVal->type = TMQ_OFFSET__LOG;
  pOffsetVal->version = ver;
}

L
Liu Jicong 已提交
194 195 196
// tqStream
int32_t tqExpandTask(STQ* pTq, SStreamTask* pTask);

H
Hongze Cheng 已提交
197 198 199 200
#ifdef __cplusplus
}
#endif

L
Liu Jicong 已提交
201
#endif /*_TD_VNODE_TQ_H_*/