tqPush.h 1.7 KB
Newer Older
L
Liu Jicong 已提交
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77
/*
 * 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/>.
 */

#ifndef _TQ_PUSH_H_
#define _TQ_PUSH_H_

#include "thash.h"
#include "trpc.h"
#include "ttimer.h"

#ifdef __cplusplus
extern "C" {
#endif

enum {
  TQ_PUSHER_TYPE__CLIENT = 1,
  TQ_PUSHER_TYPE__STREAM,
};

typedef struct {
  int8_t   type;
  int8_t   reserved[3];
  int32_t  ttl;
  int64_t  consumerId;
  SRpcMsg* pMsg;
  // SMqPollRsp* rsp;
} STqClientPusher;

typedef struct {
  int8_t  type;
  int8_t  nodeType;
  int8_t  reserved[6];
  int64_t streamId;
  SEpSet  epSet;
} STqStreamPusher;

typedef struct {
  int8_t type;  // mq or stream
} STqPusher;

typedef struct {
  SHashObj* pHash;  // <id, STqPush*>
} STqPushMgr;

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

static STqPushMgmt tqPushMgmt;

int32_t tqPushMgrInit();
void    tqPushMgrCleanUp();

STqPushMgr* tqPushMgrOpen();
void        tqPushMgrClose(STqPushMgr* pushMgr);

STqClientPusher* tqAddClientPusher(STqPushMgr* pushMgr, SRpcMsg* pMsg, int64_t consumerId, int64_t ttl);
STqStreamPusher* tqAddStreamPusher(STqPushMgr* pushMgr, int64_t streamId, SEpSet* pEpSet);

#ifdef __cplusplus
}
#endif

#endif /*_TQ_PUSH_H_*/