transSrv.c 18.6 KB
Newer Older
dengyihao's avatar
dengyihao 已提交
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17
/*
 * 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/>.
 */

#ifdef USE_UV

18
#include "transComm.h"
dengyihao's avatar
dengyihao 已提交
19 20 21 22 23 24 25 26
typedef struct SConn {
  uv_tcp_t*   pTcp;
  uv_write_t* pWriter;
  uv_timer_t* pTimer;

  uv_async_t* pWorkerAsync;
  queue       queue;
  int         ref;
dengyihao's avatar
dengyihao 已提交
27 28
  int         persist;  // persist connection or not
  SConnBuffer connBuf;  // read buf,
dengyihao's avatar
dengyihao 已提交
29
  int         count;
dengyihao's avatar
dengyihao 已提交
30 31 32
  int         inType;
  void*       pTransInst;  // rpc init
  void*       ahandle;     //
dengyihao's avatar
dengyihao 已提交
33
  void*       hostThrd;
dengyihao's avatar
dengyihao 已提交
34 35

  SRpcMsg sendMsg;
dengyihao's avatar
dengyihao 已提交
36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52
  // del later
  char secured;
  int  spi;
  char info[64];
  char user[TSDB_UNI_LEN];  // user ID for the link
  char secret[TSDB_PASSWORD_LEN];
  char ckey[TSDB_PASSWORD_LEN];  // ciphering key
} SConn;

typedef struct SWorkThrdObj {
  pthread_t       thread;
  uv_pipe_t*      pipe;
  int             fd;
  uv_loop_t*      loop;
  uv_async_t*     workerAsync;  //
  queue           conn;
  pthread_mutex_t connMtx;
dengyihao's avatar
dengyihao 已提交
53
  void*           pTransInst;
dengyihao's avatar
dengyihao 已提交
54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70
} SWorkThrdObj;

typedef struct SServerObj {
  pthread_t      thread;
  uv_tcp_t       server;
  uv_loop_t*     loop;
  int            workerIdx;
  int            numOfThreads;
  SWorkThrdObj** pThreadObj;
  uv_pipe_t**    pipe;
  uint32_t       ip;
  uint32_t       port;
} SServerObj;

static const char* notify = "a";

// refactor later
dengyihao's avatar
dengyihao 已提交
71
static int transAddAuthPart(SConn* pConn, char* msg, int msgLen);
dengyihao's avatar
dengyihao 已提交
72 73 74 75 76 77 78 79

static int uvAuthMsg(SConn* pConn, char* msg, int msgLen);

static void uvAllocConnBufferCb(uv_handle_t* handle, size_t suggested_size, uv_buf_t* buf);
static void uvAllocReadBufferCb(uv_handle_t* handle, size_t suggested_size, uv_buf_t* buf);
static void uvOnReadCb(uv_stream_t* cli, ssize_t nread, const uv_buf_t* buf);
static void uvOnTimeoutCb(uv_timer_t* handle);
static void uvOnWriteCb(uv_write_t* req, int status);
dengyihao's avatar
dengyihao 已提交
80
static void uvOnPipeWriteCb(uv_write_t* req, int status);
dengyihao's avatar
dengyihao 已提交
81 82 83 84
static void uvOnAcceptCb(uv_stream_t* stream, int status);
static void uvOnConnectionCb(uv_stream_t* q, ssize_t nread, const uv_buf_t* buf);
static void uvWorkerAsyncCb(uv_async_t* handle);

dengyihao's avatar
dengyihao 已提交
85 86
static void uvPrepareSendData(SConn* conn, uv_buf_t* wb);

87 88 89 90
// check whether already read complete packet
static bool   readComplete(SConnBuffer* buf);
static SConn* createConn();
static void   destroyConn(SConn* conn, bool clear /*clear handle or not*/);
dengyihao's avatar
dengyihao 已提交
91

92
static void uvDestroyConn(uv_handle_t* handle);
dengyihao's avatar
dengyihao 已提交
93

dengyihao's avatar
dengyihao 已提交
94
// server and worker thread
dengyihao's avatar
dengyihao 已提交
95 96 97
static void* workerThread(void* arg);
static void* acceptThread(void* arg);

dengyihao's avatar
dengyihao 已提交
98 99 100 101
// add handle loop
static bool addHandleToWorkloop(void* arg);
static bool addHandleToAcceptloop(void* arg);

dengyihao's avatar
dengyihao 已提交
102 103 104
void uvAllocReadBufferCb(uv_handle_t* handle, size_t suggested_size, uv_buf_t* buf) {
  /*
   * formate of data buffer:
dengyihao's avatar
dengyihao 已提交
105 106
   * |<--------------------------data from socket------------------------------->|
   * |<------STransMsgHead------->|<-------------------other data--------------->|
dengyihao's avatar
dengyihao 已提交
107 108 109 110 111 112
   */
  static const int CAPACITY = 1024;

  SConn*       conn = handle->data;
  SConnBuffer* pBuf = &conn->connBuf;
  if (pBuf->cap == 0) {
dengyihao's avatar
dengyihao 已提交
113
    pBuf->buf = (char*)calloc(CAPACITY, sizeof(char));
dengyihao's avatar
dengyihao 已提交
114 115 116 117
    pBuf->len = 0;
    pBuf->cap = CAPACITY;
    pBuf->left = -1;

dengyihao's avatar
dengyihao 已提交
118
    buf->base = pBuf->buf;
dengyihao's avatar
dengyihao 已提交
119 120 121 122 123
    buf->len = CAPACITY;
  } else {
    if (pBuf->len >= pBuf->cap) {
      if (pBuf->left == -1) {
        pBuf->cap *= 2;
dengyihao's avatar
dengyihao 已提交
124
        pBuf->buf = realloc(pBuf->buf, pBuf->cap);
dengyihao's avatar
dengyihao 已提交
125 126
      } else if (pBuf->len + pBuf->left > pBuf->cap) {
        pBuf->cap = pBuf->len + pBuf->left;
dengyihao's avatar
dengyihao 已提交
127
        pBuf->buf = realloc(pBuf->buf, pBuf->len + pBuf->left);
dengyihao's avatar
dengyihao 已提交
128 129
      }
    }
dengyihao's avatar
dengyihao 已提交
130
    buf->base = pBuf->buf + pBuf->len;
dengyihao's avatar
dengyihao 已提交
131 132 133 134 135 136 137 138
    buf->len = pBuf->cap - pBuf->len;
  }
}

// check data read from socket completely or not
//
static bool readComplete(SConnBuffer* data) {
  // TODO(yihao): handle pipeline later
dengyihao's avatar
dengyihao 已提交
139 140
  STransMsgHead head;
  int32_t       headLen = sizeof(head);
dengyihao's avatar
dengyihao 已提交
141
  if (data->len >= headLen) {
dengyihao's avatar
dengyihao 已提交
142 143
    memcpy((char*)&head, data->buf, headLen);
    int32_t msgLen = (int32_t)htonl((uint32_t)head.msgLen);
dengyihao's avatar
dengyihao 已提交
144 145 146
    if (msgLen > data->len) {
      data->left = msgLen - data->len;
      return false;
dengyihao's avatar
dengyihao 已提交
147
    } else if (msgLen == data->len) {
dengyihao's avatar
dengyihao 已提交
148
      return true;
dengyihao's avatar
dengyihao 已提交
149 150 151
    } else if (msgLen < data->len) {
      return false;
      // handle other packet later
dengyihao's avatar
dengyihao 已提交
152 153 154 155 156 157
    }
  } else {
    return false;
  }
}

dengyihao's avatar
dengyihao 已提交
158 159 160 161 162 163 164 165 166 167 168
// static void uvDoProcess(SRecvInfo* pRecv) {
//  // impl later
//  STransMsgHead* pHead = (STransMsgHead*)pRecv->msg;
//  SRpcInfo*      pRpc = (SRpcInfo*)pRecv->shandle;
//  SConn*         pConn = pRecv->thandle;
//  tDump(pRecv->msg, pRecv->msgLen);
//  terrno = 0;
//  // SRpcReqContext* pContest;
//
//  // do auth and check
//}
dengyihao's avatar
dengyihao 已提交
169 170

static int uvAuthMsg(SConn* pConn, char* msg, int len) {
dengyihao's avatar
dengyihao 已提交
171 172 173
  STransMsgHead* pHead = (STransMsgHead*)msg;

  int code = 0;
dengyihao's avatar
dengyihao 已提交
174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216 217 218 219 220 221 222 223 224 225

  if ((pConn->secured && pHead->spi == 0) || (pHead->spi == 0 && pConn->spi == 0)) {
    // secured link, or no authentication
    pHead->msgLen = (int32_t)htonl((uint32_t)pHead->msgLen);
    // tTrace("%s, secured link, no auth is required", pConn->info);
    return 0;
  }

  if (!rpcIsReq(pHead->msgType)) {
    // for response, if code is auth failure, it shall bypass the auth process
    code = htonl(pHead->code);
    if (code == TSDB_CODE_RPC_INVALID_TIME_STAMP || code == TSDB_CODE_RPC_AUTH_FAILURE ||
        code == TSDB_CODE_RPC_INVALID_VERSION || code == TSDB_CODE_RPC_AUTH_REQUIRED ||
        code == TSDB_CODE_MND_USER_NOT_EXIST || code == TSDB_CODE_RPC_NOT_READY) {
      pHead->msgLen = (int32_t)htonl((uint32_t)pHead->msgLen);
      // tTrace("%s, dont check authentication since code is:0x%x", pConn->info, code);
      return 0;
    }
  }

  code = 0;
  if (pHead->spi == pConn->spi) {
    // authentication
    SRpcDigest* pDigest = (SRpcDigest*)((char*)pHead + len - sizeof(SRpcDigest));

    int32_t delta;
    delta = (int32_t)htonl(pDigest->timeStamp);
    delta -= (int32_t)taosGetTimestampSec();
    if (abs(delta) > 900) {
      tWarn("%s, time diff:%d is too big, msg discarded", pConn->info, delta);
      code = TSDB_CODE_RPC_INVALID_TIME_STAMP;
    } else {
      if (rpcAuthenticateMsg(pHead, len - TSDB_AUTH_LEN, pDigest->auth, pConn->secret) < 0) {
        // tDebug("%s, authentication failed, msg discarded", pConn->info);
        code = TSDB_CODE_RPC_AUTH_FAILURE;
      } else {
        pHead->msgLen = (int32_t)htonl((uint32_t)pHead->msgLen) - sizeof(SRpcDigest);
        if (!rpcIsReq(pHead->msgType)) pConn->secured = 1;  // link is secured for client
        // tTrace("%s, message is authenticated", pConn->info);
      }
    }
  } else {
    tDebug("%s, auth spi:%d not matched with received:%d", pConn->info, pConn->spi, pHead->spi);
    code = pHead->spi ? TSDB_CODE_RPC_AUTH_FAILURE : TSDB_CODE_RPC_AUTH_REQUIRED;
  }

  return code;
}

// refers specifically to query or insert timeout
static void uvHandleActivityTimeout(uv_timer_t* handle) {
  SConn* conn = handle->data;
dengyihao's avatar
dengyihao 已提交
226
  tDebug("%p timeout since no activity", conn);
dengyihao's avatar
dengyihao 已提交
227 228 229 230 231 232
}

static void uvProcessData(SConn* pConn) {
  SRecvInfo    info;
  SRecvInfo*   p = &info;
  SConnBuffer* pBuf = &pConn->connBuf;
dengyihao's avatar
dengyihao 已提交
233
  p->msg = pBuf->buf;
dengyihao's avatar
dengyihao 已提交
234 235 236
  p->msgLen = pBuf->len;
  p->ip = 0;
  p->port = 0;
dengyihao's avatar
dengyihao 已提交
237
  p->shandle = pConn->pTransInst;  //
dengyihao's avatar
dengyihao 已提交
238 239 240
  p->thandle = pConn;
  p->chandle = NULL;

dengyihao's avatar
dengyihao 已提交
241
  STransMsgHead* pHead = (STransMsgHead*)p->msg;
dengyihao's avatar
dengyihao 已提交
242 243

  pConn->inType = pHead->msgType;
dengyihao's avatar
dengyihao 已提交
244
  assert(transIsReq(pHead->msgType));
dengyihao's avatar
dengyihao 已提交
245 246 247

  SRpcInfo* pRpc = (SRpcInfo*)p->shandle;
  // auth here
dengyihao's avatar
dengyihao 已提交
248
  // auth should not do in rpc thread
dengyihao's avatar
dengyihao 已提交
249

dengyihao's avatar
dengyihao 已提交
250 251 252 253 254
  // int8_t code = uvAuthMsg(pConn, (char*)pHead, p->msgLen);
  // if (code != 0) {
  //  terrno = code;
  //  return;
  //}
dengyihao's avatar
dengyihao 已提交
255 256
  pHead->code = htonl(pHead->code);

dengyihao's avatar
dengyihao 已提交
257
  int32_t dlen = 0;
dengyihao's avatar
dengyihao 已提交
258
  SRpcMsg rpcMsg;
dengyihao's avatar
dengyihao 已提交
259 260 261 262
  if (transDecompressMsg(NULL, 0, NULL)) {
    // add compress later
    // pHead = rpcDecompressRpcMsg(pHead);
  } else {
dengyihao's avatar
dengyihao 已提交
263
    pHead->msgLen = htonl(pHead->msgLen);
dengyihao's avatar
dengyihao 已提交
264
    // impl later
dengyihao's avatar
dengyihao 已提交
265
    //
dengyihao's avatar
dengyihao 已提交
266
  }
dengyihao's avatar
dengyihao 已提交
267
  rpcMsg.contLen = transContLenFromMsg(pHead->msgLen);
dengyihao's avatar
dengyihao 已提交
268 269 270
  rpcMsg.pCont = pHead->content;
  rpcMsg.msgType = pHead->msgType;
  rpcMsg.code = pHead->code;
dengyihao's avatar
dengyihao 已提交
271
  rpcMsg.ahandle = NULL;
dengyihao's avatar
dengyihao 已提交
272 273 274
  rpcMsg.handle = pConn;

  (*(pRpc->cfp))(pRpc->parent, &rpcMsg, NULL);
dengyihao's avatar
dengyihao 已提交
275
  // uv_timer_start(pConn->pTimer, uvHandleActivityTimeout, pRpc->idleTime * 10000, 0);
dengyihao's avatar
dengyihao 已提交
276 277 278 279 280 281
  // auth
  // validate msg type
}

void uvOnReadCb(uv_stream_t* cli, ssize_t nread, const uv_buf_t* buf) {
  // opt
dengyihao's avatar
dengyihao 已提交
282 283
  SConn*       conn = cli->data;
  SConnBuffer* pBuf = &conn->connBuf;
dengyihao's avatar
dengyihao 已提交
284 285
  if (nread > 0) {
    pBuf->len += nread;
286
    tDebug("on read %p, total read: %d, current read: %d", cli, pBuf->len, (int)nread);
dengyihao's avatar
dengyihao 已提交
287 288
    if (readComplete(pBuf)) {
      tDebug("alread read complete packet");
dengyihao's avatar
dengyihao 已提交
289
      uvProcessData(conn);
dengyihao's avatar
dengyihao 已提交
290 291 292 293 294
    } else {
      tDebug("read half packet, continue to read");
    }
    return;
  }
295 296 297
  if (nread == 0) {
    return;
  }
dengyihao's avatar
dengyihao 已提交
298
  if (nread != UV_EOF) {
299
    tDebug("read error %s", uv_err_name(nread));
dengyihao's avatar
dengyihao 已提交
300 301 302 303 304 305 306 307 308 309 310 311 312 313
  }
}
void uvAllocConnBufferCb(uv_handle_t* handle, size_t suggested_size, uv_buf_t* buf) {
  buf->base = malloc(sizeof(char));
  buf->len = 2;
}

void uvOnTimeoutCb(uv_timer_t* handle) {
  // opt
  tDebug("time out");
}

void uvOnWriteCb(uv_write_t* req, int status) {
  SConn* conn = req->data;
dengyihao's avatar
dengyihao 已提交
314 315 316 317 318

  SConnBuffer* buf = &conn->connBuf;
  buf->len = 0;
  memset(buf->buf, 0, buf->cap);
  buf->left = -1;
dengyihao's avatar
dengyihao 已提交
319 320 321
  if (status == 0) {
    tDebug("data already was written on stream");
  } else {
dengyihao's avatar
dengyihao 已提交
322
    tDebug("failed to write data, %s", uv_err_name(status));
323
    destroyConn(conn, true);
dengyihao's avatar
dengyihao 已提交
324 325 326
  }
  // opt
}
dengyihao's avatar
dengyihao 已提交
327 328 329 330 331 332 333
static void uvOnPipeWriteCb(uv_write_t* req, int status) {
  if (status == 0) {
    tDebug("success to dispatch conn to work thread");
  } else {
    tError("fail to dispatch conn to work thread");
  }
}
dengyihao's avatar
dengyihao 已提交
334

dengyihao's avatar
dengyihao 已提交
335
static void uvPrepareSendData(SConn* conn, uv_buf_t* wb) {
336 337
  // impl later;
  tDebug("prepare to send back");
dengyihao's avatar
dengyihao 已提交
338 339 340 341 342 343 344 345 346 347 348 349 350 351 352 353 354
  SRpcMsg* pMsg = &conn->sendMsg;
  if (pMsg->pCont == 0) {
    pMsg->pCont = (void*)rpcMallocCont(0);
    pMsg->contLen = 0;
  }
  STransMsgHead* pHead = transHeadFromCont(pMsg->pCont);
  pHead->msgType = conn->inType + 1;
  // add more info
  char*   msg = (char*)pHead;
  int32_t len = transMsgLenFromCont(pMsg->contLen);
  if (transCompressMsg(msg, len, NULL)) {
    // impl later
  }
  pHead->msgLen = htonl(len);
  wb->base = msg;
  wb->len = len;
}
dengyihao's avatar
dengyihao 已提交
355
void uvWorkerAsyncCb(uv_async_t* handle) {
dengyihao's avatar
dengyihao 已提交
356
  SWorkThrdObj* pThrd = handle->data;
dengyihao's avatar
dengyihao 已提交
357
  SConn*        conn = NULL;
dengyihao's avatar
dengyihao 已提交
358 359
  queue         wq;
  // batch process to avoid to lock/unlock frequently
dengyihao's avatar
dengyihao 已提交
360
  pthread_mutex_lock(&pThrd->connMtx);
dengyihao's avatar
dengyihao 已提交
361
  QUEUE_MOVE(&pThrd->conn, &wq);
dengyihao's avatar
dengyihao 已提交
362
  pthread_mutex_unlock(&pThrd->connMtx);
dengyihao's avatar
dengyihao 已提交
363 364 365 366 367 368 369 370 371

  while (!QUEUE_IS_EMPTY(&wq)) {
    queue* head = QUEUE_HEAD(&wq);
    QUEUE_REMOVE(head);
    SConn* conn = QUEUE_DATA(head, SConn, queue);
    if (conn == NULL) {
      tError("except occurred, do nothing");
      return;
    }
dengyihao's avatar
dengyihao 已提交
372 373
    uv_buf_t wb;
    uvPrepareSendData(conn, &wb);
dengyihao's avatar
dengyihao 已提交
374 375
    uv_timer_stop(conn->pTimer);

dengyihao's avatar
dengyihao 已提交
376
    uv_write(conn->pWriter, (uv_stream_t*)conn->pTcp, &wb, 1, uvOnWriteCb);
dengyihao's avatar
dengyihao 已提交
377 378 379 380 381 382 383 384 385 386 387 388 389 390 391 392 393 394
  }
}

void uvOnAcceptCb(uv_stream_t* stream, int status) {
  if (status == -1) {
    return;
  }
  SServerObj* pObj = container_of(stream, SServerObj, server);

  uv_tcp_t* cli = (uv_tcp_t*)malloc(sizeof(uv_tcp_t));
  uv_tcp_init(pObj->loop, cli);

  if (uv_accept(stream, (uv_stream_t*)cli) == 0) {
    uv_write_t* wr = (uv_write_t*)malloc(sizeof(uv_write_t));

    uv_buf_t buf = uv_buf_init((char*)notify, strlen(notify));

    pObj->workerIdx = (pObj->workerIdx + 1) % pObj->numOfThreads;
dengyihao's avatar
dengyihao 已提交
395

dengyihao's avatar
dengyihao 已提交
396
    tDebug("new conntion accepted by main server, dispatch to %dth worker-thread", pObj->workerIdx);
dengyihao's avatar
dengyihao 已提交
397
    uv_write2(wr, (uv_stream_t*)&(pObj->pipe[pObj->workerIdx][0]), &buf, 1, (uv_stream_t*)cli, uvOnPipeWriteCb);
dengyihao's avatar
dengyihao 已提交
398 399
  } else {
    uv_close((uv_handle_t*)cli, NULL);
400
    free(cli);
dengyihao's avatar
dengyihao 已提交
401 402 403 404 405 406 407 408 409
  }
}
void uvOnConnectionCb(uv_stream_t* q, ssize_t nread, const uv_buf_t* buf) {
  tDebug("connection coming");
  if (nread < 0) {
    if (nread != UV_EOF) {
      tError("read error %s", uv_err_name(nread));
    }
    // TODO(log other failure reason)
410
    // uv_close((uv_handle_t*)q, NULL);
dengyihao's avatar
dengyihao 已提交
411 412 413 414 415 416 417 418 419 420 421 422 423 424 425 426 427 428
    return;
  }
  // free memory allocated by
  assert(nread == strlen(notify));
  assert(buf->base[0] == notify[0]);
  free(buf->base);

  SWorkThrdObj* pThrd = q->data;

  uv_pipe_t* pipe = (uv_pipe_t*)q;
  if (!uv_pipe_pending_count(pipe)) {
    tError("No pending count");
    return;
  }

  uv_handle_type pending = uv_pipe_pending_type(pipe);
  assert(pending == UV_TCP);

429
  SConn* pConn = createConn();
dengyihao's avatar
dengyihao 已提交
430
  pConn->pTransInst = pThrd->pTransInst;
dengyihao's avatar
dengyihao 已提交
431 432 433 434 435 436 437 438 439 440 441 442 443 444 445 446 447 448 449 450 451 452 453
  /* init conn timer*/
  pConn->pTimer = malloc(sizeof(uv_timer_t));
  uv_timer_init(pThrd->loop, pConn->pTimer);
  pConn->pTimer->data = pConn;

  pConn->hostThrd = pThrd;
  pConn->pWorkerAsync = pThrd->workerAsync;  // thread safty

  // init client handle
  pConn->pTcp = (uv_tcp_t*)malloc(sizeof(uv_tcp_t));
  uv_tcp_init(pThrd->loop, pConn->pTcp);
  pConn->pTcp->data = pConn;

  // init write request, just
  pConn->pWriter = calloc(1, sizeof(uv_write_t));
  pConn->pWriter->data = pConn;

  if (uv_accept(q, (uv_stream_t*)(pConn->pTcp)) == 0) {
    uv_os_fd_t fd;
    uv_fileno((const uv_handle_t*)pConn->pTcp, &fd);
    tDebug("new connection created: %d", fd);
    uv_read_start((uv_stream_t*)(pConn->pTcp), uvAllocReadBufferCb, uvOnReadCb);
  } else {
dengyihao's avatar
dengyihao 已提交
454
    tDebug("failed to create new connection");
455
    destroyConn(pConn, true);
dengyihao's avatar
dengyihao 已提交
456 457 458 459 460 461 462 463
  }
}

void* acceptThread(void* arg) {
  // opt
  SServerObj* srv = (SServerObj*)arg;
  uv_run(srv->loop, UV_RUN_DEFAULT);
}
dengyihao's avatar
dengyihao 已提交
464 465
static bool addHandleToWorkloop(void* arg) {
  SWorkThrdObj* pThrd = arg;
dengyihao's avatar
dengyihao 已提交
466
  pThrd->loop = (uv_loop_t*)malloc(sizeof(uv_loop_t));
dengyihao's avatar
dengyihao 已提交
467 468 469
  if (0 != uv_loop_init(pThrd->loop)) {
    return false;
  }
dengyihao's avatar
dengyihao 已提交
470 471

  // SRpcInfo* pRpc = pThrd->shandle;
dengyihao's avatar
dengyihao 已提交
472
  uv_pipe_init(pThrd->loop, pThrd->pipe, 1);
dengyihao's avatar
dengyihao 已提交
473 474 475 476 477 478 479 480 481
  uv_pipe_open(pThrd->pipe, pThrd->fd);

  pThrd->pipe->data = pThrd;

  QUEUE_INIT(&pThrd->conn);
  pthread_mutex_init(&pThrd->connMtx, NULL);

  pThrd->workerAsync = malloc(sizeof(uv_async_t));
  uv_async_init(pThrd->loop, pThrd->workerAsync, uvWorkerAsyncCb);
dengyihao's avatar
dengyihao 已提交
482
  pThrd->workerAsync->data = pThrd;
dengyihao's avatar
dengyihao 已提交
483 484

  uv_read_start((uv_stream_t*)pThrd->pipe, uvAllocConnBufferCb, uvOnConnectionCb);
dengyihao's avatar
dengyihao 已提交
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
  return true;
}

static bool addHandleToAcceptloop(void* arg) {
  // impl later
  SServerObj* srv = arg;

  int err = 0;
  if ((err = uv_tcp_init(srv->loop, &srv->server)) != 0) {
    tError("failed to init accept server: %s", uv_err_name(err));
    return false;
  }

  struct sockaddr_in bind_addr;

  uv_ip4_addr("0.0.0.0", srv->port, &bind_addr);
  if ((err = uv_tcp_bind(&srv->server, (const struct sockaddr*)&bind_addr, 0)) != 0) {
    tError("failed to bind: %s", uv_err_name(err));
    return false;
  }
  if ((err = uv_listen((uv_stream_t*)&srv->server, 128, uvOnAcceptCb)) != 0) {
    tError("failed to listen: %s", uv_err_name(err));
    return false;
  }
  return true;
dengyihao's avatar
dengyihao 已提交
510 511 512
}
void* workerThread(void* arg) {
  SWorkThrdObj* pThrd = (SWorkThrdObj*)arg;
dengyihao's avatar
dengyihao 已提交
513 514 515
  uv_run(pThrd->loop, UV_RUN_DEFAULT);
}

516
static SConn* createConn() {
dengyihao's avatar
dengyihao 已提交
517 518 519
  SConn* pConn = (SConn*)calloc(1, sizeof(SConn));
  return pConn;
}
dengyihao's avatar
dengyihao 已提交
520 521 522 523
static void connCloseCb(uv_handle_t* handle) {
  // impl later
  //
}
524
static void destroyConn(SConn* conn, bool clear) {
dengyihao's avatar
dengyihao 已提交
525 526 527
  if (conn == NULL) {
    return;
  }
528 529 530 531
  if (clear) {
    uv_handle_t handle = *((uv_handle_t*)conn->pTcp);
    uv_close(&handle, NULL);
  }
dengyihao's avatar
dengyihao 已提交
532 533 534
  uv_timer_stop(conn->pTimer);
  free(conn->pTimer);
  free(conn->pTcp);
dengyihao's avatar
dengyihao 已提交
535
  free(conn->connBuf.buf);
dengyihao's avatar
dengyihao 已提交
536
  free(conn->pWriter);
537
  free(conn);
dengyihao's avatar
dengyihao 已提交
538
}
539
static void uvDestroyConn(uv_handle_t* handle) {
dengyihao's avatar
dengyihao 已提交
540
  SConn* conn = handle->data;
541
  destroyConn(conn, false);
dengyihao's avatar
dengyihao 已提交
542
}
dengyihao's avatar
dengyihao 已提交
543 544
static int transAddAuthPart(SConn* pConn, char* msg, int msgLen) {
  STransMsgHead* pHead = (STransMsgHead*)msg;
dengyihao's avatar
dengyihao 已提交
545 546 547 548

  if (pConn->spi && pConn->secured == 0) {
    // add auth part
    pHead->spi = pConn->spi;
dengyihao's avatar
dengyihao 已提交
549
    STransDigestMsg* pDigest = (STransDigestMsg*)(msg + msgLen);
dengyihao's avatar
dengyihao 已提交
550 551 552
    pDigest->timeStamp = htonl(taosGetTimestampSec());
    msgLen += sizeof(SRpcDigest);
    pHead->msgLen = (int32_t)htonl((uint32_t)msgLen);
dengyihao's avatar
dengyihao 已提交
553 554
    // transBuildAuthHead(pHead, msgLen - TSDB_AUTH_LEN, pDigest->auth, pConn->secret);
    // transBuildAuthHead(pHead, msgLen - TSDB_AUTH_LEN, pDigest->auth, pConn->secret);
dengyihao's avatar
dengyihao 已提交
555 556 557 558 559 560 561 562 563 564 565 566 567 568 569 570 571 572 573 574 575
  } else {
    pHead->spi = 0;
    pHead->msgLen = (int32_t)htonl((uint32_t)msgLen);
  }

  return msgLen;
}

void* taosInitServer(uint32_t ip, uint32_t port, char* label, int numOfThreads, void* fp, void* shandle) {
  SServerObj* srv = calloc(1, sizeof(SServerObj));
  srv->loop = (uv_loop_t*)malloc(sizeof(uv_loop_t));
  srv->numOfThreads = numOfThreads;
  srv->workerIdx = 0;
  srv->pThreadObj = (SWorkThrdObj**)calloc(srv->numOfThreads, sizeof(SWorkThrdObj*));
  srv->pipe = (uv_pipe_t**)calloc(srv->numOfThreads, sizeof(uv_pipe_t*));
  srv->ip = ip;
  srv->port = port;
  uv_loop_init(srv->loop);

  for (int i = 0; i < srv->numOfThreads; i++) {
    SWorkThrdObj* thrd = (SWorkThrdObj*)calloc(1, sizeof(SWorkThrdObj));
dengyihao's avatar
dengyihao 已提交
576
    srv->pThreadObj[i] = thrd;
dengyihao's avatar
dengyihao 已提交
577

dengyihao's avatar
dengyihao 已提交
578 579 580
    srv->pipe[i] = (uv_pipe_t*)calloc(2, sizeof(uv_pipe_t));
    int fds[2];
    if (uv_socketpair(AF_UNIX, SOCK_STREAM, fds, UV_NONBLOCK_PIPE, UV_NONBLOCK_PIPE) != 0) {
dengyihao's avatar
dengyihao 已提交
581
      goto End;
dengyihao's avatar
dengyihao 已提交
582 583 584 585
    }
    uv_pipe_init(srv->loop, &(srv->pipe[i][0]), 1);
    uv_pipe_open(&(srv->pipe[i][0]), fds[1]);  // init write

dengyihao's avatar
dengyihao 已提交
586
    thrd->pTransInst = shandle;
dengyihao's avatar
dengyihao 已提交
587 588
    thrd->fd = fds[0];
    thrd->pipe = &(srv->pipe[i][1]);  // init read
dengyihao's avatar
dengyihao 已提交
589

dengyihao's avatar
dengyihao 已提交
590 591 592
    if (false == addHandleToWorkloop(thrd)) {
      goto End;
    }
dengyihao's avatar
dengyihao 已提交
593 594 595 596 597 598 599 600 601
    int err = pthread_create(&(thrd->thread), NULL, workerThread, (void*)(thrd));
    if (err == 0) {
      tDebug("sucess to create worker-thread %d", i);
      // printf("thread %d create\n", i);
    } else {
      // TODO: clear all other resource later
      tError("failed to create worker-thread %d", i);
    }
  }
dengyihao's avatar
dengyihao 已提交
602 603 604
  if (false == addHandleToAcceptloop(srv)) {
    goto End;
  }
dengyihao's avatar
dengyihao 已提交
605 606 607 608 609 610 611 612
  int err = pthread_create(&srv->thread, NULL, acceptThread, (void*)srv);
  if (err == 0) {
    tDebug("success to create accept-thread");
  } else {
    // clear all resource later
  }

  return srv;
dengyihao's avatar
dengyihao 已提交
613 614 615 616 617 618 619 620 621 622 623 624 625
End:
  taosCloseServer(srv);
  return NULL;
}

void destroyWorkThrd(SWorkThrdObj* pThrd) {
  if (pThrd == NULL) {
    return;
  }
  pthread_join(pThrd->thread, NULL);
  // free(srv->pipe[i]);
  free(pThrd->loop);
  free(pThrd);
dengyihao's avatar
dengyihao 已提交
626
}
dengyihao's avatar
dengyihao 已提交
627 628 629
void taosCloseServer(void* arg) {
  // impl later
  SServerObj* srv = arg;
dengyihao's avatar
dengyihao 已提交
630
  for (int i = 0; i < srv->numOfThreads; i++) {
dengyihao's avatar
dengyihao 已提交
631
    destroyWorkThrd(srv->pThreadObj[i]);
dengyihao's avatar
dengyihao 已提交
632 633 634 635 636
  }
  free(srv->loop);
  free(srv->pipe);
  free(srv->pThreadObj);
  pthread_join(srv->thread, NULL);
dengyihao's avatar
dengyihao 已提交
637
  free(srv);
dengyihao's avatar
dengyihao 已提交
638
}
dengyihao's avatar
dengyihao 已提交
639 640 641 642 643 644

void rpcSendResponse(const SRpcMsg* pMsg) {
  SConn*        pConn = pMsg->handle;
  SWorkThrdObj* pThrd = pConn->hostThrd;

  // opt later
dengyihao's avatar
dengyihao 已提交
645
  pConn->sendMsg = *pMsg;
dengyihao's avatar
dengyihao 已提交
646 647 648 649 650 651 652 653
  pthread_mutex_lock(&pThrd->connMtx);
  QUEUE_PUSH(&pThrd->conn, &pConn->queue);
  pthread_mutex_unlock(&pThrd->connMtx);

  uv_async_send(pConn->pWorkerAsync);
}

#endif