rpcTcp.c 19.5 KB
Newer Older
H
hzcheng 已提交
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/>.
 */

dengyihao's avatar
dengyihao 已提交
16
#include "rpcTcp.h"
dengyihao's avatar
dengyihao 已提交
17
#ifdef USE_UV
dengyihao's avatar
dengyihao 已提交
18
#include <uv.h>
dengyihao's avatar
dengyihao 已提交
19
#endif
S
slguan 已提交
20
#include "os.h"
dengyihao's avatar
dengyihao 已提交
21 22
#include "rpcHead.h"
#include "rpcLog.h"
23
#include "taosdef.h"
F
freemine 已提交
24
#include "taoserror.h"
dengyihao's avatar
dengyihao 已提交
25
#include "tutil.h"
H
hzcheng 已提交
26

dengyihao's avatar
dengyihao 已提交
27 28 29
#ifdef USE_UV

#else
J
Jeff Tao 已提交
30
typedef struct SFdObj {
dengyihao's avatar
dengyihao 已提交
31 32 33
  void *             signature;
  SOCKET             fd;       // TCP socket FD
  void *             thandle;  // handle from upper layer, like TAOS
J
Jeff Tao 已提交
34 35
  uint32_t           ip;
  uint16_t           port;
dengyihao's avatar
dengyihao 已提交
36
  int16_t            closedByApp;  // 1: already closed by App
J
Jeff Tao 已提交
37
  struct SThreadObj *pThreadObj;
dengyihao's avatar
dengyihao 已提交
38 39
  struct SFdObj *    prev;
  struct SFdObj *    next;
H
hzcheng 已提交
40 41
} SFdObj;

J
Jeff Tao 已提交
42
typedef struct SThreadObj {
H
hzcheng 已提交
43 44
  pthread_t       thread;
  SFdObj *        pHead;
J
Jeff Tao 已提交
45
  pthread_mutex_t mutex;
J
jtao1735 已提交
46
  uint32_t        ip;
47
  bool            stop;
48
  EpollFd         pollFd;
H
hzcheng 已提交
49 50
  int             numOfFds;
  int             threadId;
51
  char            label[TSDB_LABEL_LEN];
dengyihao's avatar
dengyihao 已提交
52 53
  void *          shandle;  // handle passed by upper layer during server initialization
  void *(*processData)(SRecvInfo *pPacket);
H
hzcheng 已提交
54 55
} SThreadObj;

Y
yihaoDeng 已提交
56
typedef struct {
dengyihao's avatar
dengyihao 已提交
57 58 59
  char         label[TSDB_LABEL_LEN];
  int32_t      index;
  int          numOfThreads;
Y
yihaoDeng 已提交
60 61 62
  SThreadObj **pThreadObj;
} SClientObj;

H
hzcheng 已提交
63
typedef struct {
dengyihao's avatar
dengyihao 已提交
64 65 66 67 68 69 70 71
  SOCKET       fd;
  uint32_t     ip;
  uint16_t     port;
  int8_t       stop;
  int8_t       reserve;
  char         label[TSDB_LABEL_LEN];
  int          numOfThreads;
  void *       shandle;
陶建辉(Jeff)'s avatar
TD-1645  
陶建辉(Jeff) 已提交
72
  SThreadObj **pThreadObj;
dengyihao's avatar
dengyihao 已提交
73
  pthread_t    thread;
H
hzcheng 已提交
74 75
} SServerObj;

dengyihao's avatar
dengyihao 已提交
76
static void *  taosProcessTcpData(void *param);
S
TD-1057  
Shengliang Guan 已提交
77
static SFdObj *taosMallocFdObj(SThreadObj *pThreadObj, SOCKET fd);
J
Jeff Tao 已提交
78 79
static void    taosFreeFdObj(SFdObj *pFdObj);
static void    taosReportBrokenLink(SFdObj *pFdObj);
dengyihao's avatar
dengyihao 已提交
80
static void *  taosAcceptTcpConnection(void *arg);
H
hzcheng 已提交
81

J
jtao1735 已提交
82
void *taosInitTcpServer(uint32_t ip, uint16_t port, char *label, int numOfThreads, void *fp, void *shandle) {
J
Jeff Tao 已提交
83 84
  SServerObj *pServerObj;
  SThreadObj *pThreadObj;
H
hzcheng 已提交
85

J
Jeff Tao 已提交
86
  pServerObj = (SServerObj *)calloc(sizeof(SServerObj), 1);
陶建辉(Jeff)'s avatar
陶建辉(Jeff) 已提交
87 88
  if (pServerObj == NULL) {
    tError("TCP:%s no enough memory", label);
F
freemine 已提交
89
    terrno = TAOS_SYSTEM_ERROR(errno);
陶建辉(Jeff)'s avatar
陶建辉(Jeff) 已提交
90 91 92
    return NULL;
  }

93
  pServerObj->fd = -1;
94
  taosResetPthread(&pServerObj->thread);
J
jtao1735 已提交
95
  pServerObj->ip = ip;
96
  pServerObj->port = port;
陶建辉(Jeff)'s avatar
陶建辉(Jeff) 已提交
97
  tstrncpy(pServerObj->label, label, sizeof(pServerObj->label));
98 99
  pServerObj->numOfThreads = numOfThreads;

陶建辉(Jeff)'s avatar
TD-1645  
陶建辉(Jeff) 已提交
100
  pServerObj->pThreadObj = (SThreadObj **)calloc(sizeof(SThreadObj *), numOfThreads);
101 102
  if (pServerObj->pThreadObj == NULL) {
    tError("TCP:%s no enough memory", label);
F
freemine 已提交
103
    terrno = TAOS_SYSTEM_ERROR(errno);
J
Jeff Tao 已提交
104
    free(pServerObj);
105
    return NULL;
H
hzcheng 已提交
106
  }
J
Jeff Tao 已提交
107

dengyihao's avatar
dengyihao 已提交
108
  int            code = 0;
109 110 111 112
  pthread_attr_t thattr;
  pthread_attr_init(&thattr);
  pthread_attr_setdetachstate(&thattr, PTHREAD_CREATE_JOINABLE);

F
freemine 已提交
113
  // initialize parameters in case it may encounter error later
J
Jeff Tao 已提交
114
  for (int i = 0; i < numOfThreads; ++i) {
陶建辉(Jeff)'s avatar
TD-1645  
陶建辉(Jeff) 已提交
115 116 117
    pThreadObj = (SThreadObj *)calloc(sizeof(SThreadObj), 1);
    if (pThreadObj == NULL) {
      tError("TCP:%s no enough memory", label);
F
freemine 已提交
118
      terrno = TAOS_SYSTEM_ERROR(errno);
dengyihao's avatar
dengyihao 已提交
119
      for (int j = 0; j < i; ++j) free(pServerObj->pThreadObj[j]);
陶建辉(Jeff)'s avatar
TD-1645  
陶建辉(Jeff) 已提交
120 121 122 123
      free(pServerObj->pThreadObj);
      free(pServerObj);
      return NULL;
    }
F
freemine 已提交
124

陶建辉(Jeff)'s avatar
TD-1645  
陶建辉(Jeff) 已提交
125
    pServerObj->pThreadObj[i] = pThreadObj;
126
    pThreadObj->pollFd = -1;
127
    taosResetPthread(&pThreadObj->thread);
128
    pThreadObj->processData = fp;
陶建辉(Jeff)'s avatar
陶建辉(Jeff) 已提交
129
    tstrncpy(pThreadObj->label, label, sizeof(pThreadObj->label));
130
    pThreadObj->shandle = shandle;
Y
yihaoDeng 已提交
131
    pThreadObj->stop = false;
132
  }
H
hzcheng 已提交
133

134 135
  // initialize mutex, thread, fd which may fail
  for (int i = 0; i < numOfThreads; ++i) {
陶建辉(Jeff)'s avatar
TD-1645  
陶建辉(Jeff) 已提交
136
    pThreadObj = pServerObj->pThreadObj[i];
J
Jeff Tao 已提交
137
    code = pthread_mutex_init(&(pThreadObj->mutex), NULL);
J
Jeff Tao 已提交
138 139
    if (code < 0) {
      tError("%s failed to init TCP process data mutex(%s)", label, strerror(errno));
140
      break;
141
    }
H
hzcheng 已提交
142

143
    pThreadObj->pollFd = (EpollFd)epoll_create(10);  // size does not matter
144 145
    if (pThreadObj->pollFd < 0) {
      tError("%s failed to create TCP epoll", label);
J
Jeff Tao 已提交
146 147
      code = -1;
      break;
148
    }
H
hzcheng 已提交
149

J
Jeff Tao 已提交
150
    code = pthread_create(&(pThreadObj->thread), &thattr, taosProcessTcpData, (void *)(pThreadObj));
J
Jeff Tao 已提交
151 152 153
    if (code != 0) {
      tError("%s failed to create TCP process data thread(%s)", label, strerror(errno));
      break;
154
    }
H
hzcheng 已提交
155

156
    pThreadObj->threadId = i;
H
hzcheng 已提交
157 158
  }

159
  pServerObj->fd = taosOpenTcpServerSocket(pServerObj->ip, pServerObj->port);
F
freemine 已提交
160
  if (pServerObj->fd < 0) code = -1;
161

F
freemine 已提交
162
  if (code == 0) {
163
    code = pthread_create(&pServerObj->thread, &thattr, taosAcceptTcpConnection, (void *)pServerObj);
J
Jeff Tao 已提交
164
    if (code != 0) {
S
TD-2805  
Shengliang Guan 已提交
165
      tError("%s failed to create TCP accept thread(%s)", label, strerror(code));
J
Jeff Tao 已提交
166
    }
H
hzcheng 已提交
167 168
  }

J
Jeff Tao 已提交
169
  if (code != 0) {
F
freemine 已提交
170
    terrno = TAOS_SYSTEM_ERROR(errno);
171
    taosCleanUpTcpServer(pServerObj);
J
Jeff Tao 已提交
172 173
    pServerObj = NULL;
  } else {
174
    tDebug("%s TCP server is initialized, ip:0x%x port:%hu numOfThreads:%d", label, ip, port, numOfThreads);
J
Jeff Tao 已提交
175
  }
H
hzcheng 已提交
176

177
  pthread_attr_destroy(&thattr);
178
  return (void *)pServerObj;
H
hzcheng 已提交
179 180
}

dengyihao's avatar
dengyihao 已提交
181 182 183 184
static void taosStopTcpThread(SThreadObj *pThreadObj) {
  if (pThreadObj == NULL) {
    return;
  }
F
freemine 已提交
185 186
  // save thread into local variable and signal thread to stop
  pthread_t thread = pThreadObj->thread;
187 188 189
  if (!taosCheckPthreadValid(thread)) {
    return;
  }
dengyihao's avatar
dengyihao 已提交
190
  pThreadObj->stop = true;
dengyihao's avatar
dengyihao 已提交
191
  if (taosComparePthread(thread, pthread_self())) {
陶建辉(Jeff)'s avatar
TD-1645  
陶建辉(Jeff) 已提交
192 193 194
    pthread_detach(pthread_self());
    return;
  }
dengyihao's avatar
dengyihao 已提交
195
  pthread_join(thread, NULL);
196 197
}

198
void taosStopTcpServer(void *handle) {
199
  SServerObj *pServerObj = handle;
H
hzcheng 已提交
200 201

  if (pServerObj == NULL) return;
S
TD-1207  
Shengliang Guan 已提交
202
  pServerObj->stop = 1;
陶建辉(Jeff)'s avatar
TD-1645  
陶建辉(Jeff) 已提交
203

S
TD-1207  
Shengliang Guan 已提交
204
  if (pServerObj->fd >= 0) {
dengyihao's avatar
dengyihao 已提交
205
    taosShutDownSocketRD(pServerObj->fd);
S
TD-1207  
Shengliang Guan 已提交
206
  }
陶建辉(Jeff)'s avatar
TD-1645  
陶建辉(Jeff) 已提交
207
  if (taosCheckPthreadValid(pServerObj->thread)) {
208
    if (taosComparePthread(pServerObj->thread, pthread_self())) {
陶建辉(Jeff)'s avatar
TD-1645  
陶建辉(Jeff) 已提交
209 210 211 212 213
      pthread_detach(pthread_self());
    } else {
      pthread_join(pServerObj->thread, NULL);
    }
  }
陶建辉(Jeff)'s avatar
陶建辉(Jeff) 已提交
214

215
  tDebug("%s TCP server is stopped", pServerObj->label);
216 217 218 219 220 221
}

void taosCleanUpTcpServer(void *handle) {
  SServerObj *pServerObj = handle;
  SThreadObj *pThreadObj;
  if (pServerObj == NULL) return;
H
hzcheng 已提交
222

J
Jeff Tao 已提交
223
  for (int i = 0; i < pServerObj->numOfThreads; ++i) {
陶建辉(Jeff)'s avatar
TD-1645  
陶建辉(Jeff) 已提交
224
    pThreadObj = pServerObj->pThreadObj[i];
225
    taosStopTcpThread(pThreadObj);
H
hzcheng 已提交
226 227
  }

228
  tDebug("%s TCP server is cleaned up", pServerObj->label);
H
hzcheng 已提交
229

S
TD-1848  
Shengliang Guan 已提交
230 231
  tfree(pServerObj->pThreadObj);
  tfree(pServerObj);
H
hzcheng 已提交
232 233
}

234
static void *taosAcceptTcpConnection(void *arg) {
S
TD-1057  
Shengliang Guan 已提交
235
  SOCKET             connFd = -1;
J
Jeff Tao 已提交
236 237
  struct sockaddr_in caddr;
  int                threadId = 0;
dengyihao's avatar
dengyihao 已提交
238 239
  SThreadObj *       pThreadObj;
  SServerObj *       pServerObj;
J
Jeff Tao 已提交
240 241

  pServerObj = (SServerObj *)arg;
242
  tDebug("%s TCP server is ready, ip:0x%x:%hu", pServerObj->label, pServerObj->ip, pServerObj->port);
243
  setThreadName("acceptTcpConn");
J
Jeff Tao 已提交
244 245 246

  while (1) {
    socklen_t addrlen = sizeof(caddr);
247
    connFd = accept(pServerObj->fd, (struct sockaddr *)&caddr, &addrlen);
S
TD-1207  
Shengliang Guan 已提交
248 249 250 251 252
    if (pServerObj->stop) {
      tDebug("%s TCP server stop accepting new connections", pServerObj->label);
      break;
    }

253 254
    if (connFd == -1) {
      if (errno == EINVAL) {
255
        tDebug("%s TCP server stop accepting new connections, exiting", pServerObj->label);
256 257
        break;
      }
陶建辉(Jeff)'s avatar
陶建辉(Jeff) 已提交
258

259
      tError("%s TCP accept failure(%s)", pServerObj->label, strerror(errno));
J
Jeff Tao 已提交
260 261 262 263
      continue;
    }

    taosKeepTcpAlive(connFd);
dengyihao's avatar
dengyihao 已提交
264 265
    struct timeval to = {5, 0};
    int32_t        ret = taosSetSockOpt(connFd, SOL_SOCKET, SO_RCVTIMEO, &to, sizeof(to));
Y
yihaoDeng 已提交
266 267
    if (ret != 0) {
      taosCloseSocket(connFd);
268 269
      tError("%s failed to set recv timeout fd(%s)for connection from:%s:%hu", pServerObj->label, strerror(errno),
             taosInetNtoa(caddr.sin_addr), htons(caddr.sin_port));
Y
yihaoDeng 已提交
270 271
      continue;
    }
F
freemine 已提交
272

J
Jeff Tao 已提交
273
    // pick up the thread to handle this connection
陶建辉(Jeff)'s avatar
TD-1645  
陶建辉(Jeff) 已提交
274
    pThreadObj = pServerObj->pThreadObj[threadId];
J
Jeff Tao 已提交
275 276 277 278

    SFdObj *pFdObj = taosMallocFdObj(pThreadObj, connFd);
    if (pFdObj) {
      pFdObj->ip = caddr.sin_addr.s_addr;
陶建辉(Jeff)'s avatar
陶建辉(Jeff) 已提交
279
      pFdObj->port = htons(caddr.sin_port);
280
      tDebug("%s new TCP connection from %s:%hu, fd:%d FD:%p numOfFds:%d", pServerObj->label,
dengyihao's avatar
dengyihao 已提交
281
             taosInetNtoa(caddr.sin_addr), pFdObj->port, connFd, pFdObj, pThreadObj->numOfFds);
J
Jeff Tao 已提交
282
    } else {
S
TD-1057  
Shengliang Guan 已提交
283
      taosCloseSocket(connFd);
284 285
      tError("%s failed to malloc FdObj(%s) for connection from:%s:%hu", pServerObj->label, strerror(errno),
             taosInetNtoa(caddr.sin_addr), htons(caddr.sin_port));
F
freemine 已提交
286
    }
J
Jeff Tao 已提交
287 288 289 290 291

    // pick up next thread for next connection
    threadId++;
    threadId = threadId % pServerObj->numOfThreads;
  }
292

S
TD-1057  
Shengliang Guan 已提交
293
  taosCloseSocket(pServerObj->fd);
294
  return NULL;
J
Jeff Tao 已提交
295 296
}

Y
yihaoDeng 已提交
297 298 299 300
void *taosInitTcpClient(uint32_t ip, uint16_t port, char *label, int numOfThreads, void *fp, void *shandle) {
  SClientObj *pClientObj = (SClientObj *)calloc(1, sizeof(SClientObj));
  if (pClientObj == NULL) {
    tError("TCP:%s no enough memory", label);
F
freemine 已提交
301
    terrno = TAOS_SYSTEM_ERROR(errno);
J
Jeff Tao 已提交
302
    return NULL;
F
freemine 已提交
303
  }
J
Jeff Tao 已提交
304

Y
yihaoDeng 已提交
305 306
  tstrncpy(pClientObj->label, label, sizeof(pClientObj->label));
  pClientObj->numOfThreads = numOfThreads;
dengyihao's avatar
dengyihao 已提交
307
  pClientObj->pThreadObj = (SThreadObj **)calloc(numOfThreads, sizeof(SThreadObj *));
Y
yihaoDeng 已提交
308 309 310
  if (pClientObj->pThreadObj == NULL) {
    tError("TCP:%s no enough memory", label);
    tfree(pClientObj);
F
freemine 已提交
311
    terrno = TAOS_SYSTEM_ERROR(errno);
J
Jeff Tao 已提交
312 313
  }

dengyihao's avatar
dengyihao 已提交
314
  int            code = 0;
Y
yihaoDeng 已提交
315
  pthread_attr_t thattr;
J
Jeff Tao 已提交
316 317
  pthread_attr_init(&thattr);
  pthread_attr_setdetachstate(&thattr, PTHREAD_CREATE_JOINABLE);
Y
yihaoDeng 已提交
318 319 320 321 322

  for (int i = 0; i < numOfThreads; ++i) {
    SThreadObj *pThreadObj = (SThreadObj *)calloc(1, sizeof(SThreadObj));
    if (pThreadObj == NULL) {
      tError("TCP:%s no enough memory", label);
F
freemine 已提交
323
      terrno = TAOS_SYSTEM_ERROR(errno);
dengyihao's avatar
dengyihao 已提交
324
      for (int j = 0; j < i; ++j) free(pClientObj->pThreadObj[j]);
Y
yihaoDeng 已提交
325 326 327 328 329 330
      free(pClientObj);
      pthread_attr_destroy(&thattr);
      return NULL;
    }
    pClientObj->pThreadObj[i] = pThreadObj;
    taosResetPthread(&pThreadObj->thread);
dengyihao's avatar
dengyihao 已提交
331 332
    pThreadObj->ip = ip;
    pThreadObj->stop = false;
Y
yihaoDeng 已提交
333 334 335
    tstrncpy(pThreadObj->label, label, sizeof(pThreadObj->label));
    pThreadObj->shandle = shandle;
    pThreadObj->processData = fp;
J
Jeff Tao 已提交
336 337
  }

Y
yihaoDeng 已提交
338 339 340 341 342 343 344 345 346 347 348 349 350 351 352
  // initialize mutex, thread, fd which may fail
  for (int i = 0; i < numOfThreads; ++i) {
    SThreadObj *pThreadObj = pClientObj->pThreadObj[i];
    code = pthread_mutex_init(&(pThreadObj->mutex), NULL);
    if (code < 0) {
      tError("%s failed to init TCP process data mutex(%s)", label, strerror(errno));
      break;
    }

    pThreadObj->pollFd = (int64_t)epoll_create(10);  // size does not matter
    if (pThreadObj->pollFd < 0) {
      tError("%s failed to create TCP epoll", label);
      code = -1;
      break;
    }
J
Jeff Tao 已提交
353

Y
yihaoDeng 已提交
354 355 356 357 358 359 360 361
    code = pthread_create(&(pThreadObj->thread), &thattr, taosProcessTcpData, (void *)(pThreadObj));
    if (code != 0) {
      tError("%s failed to create TCP process data thread(%s)", label, strerror(errno));
      break;
    }
    pThreadObj->threadId = i;
  }
  if (code != 0) {
F
freemine 已提交
362
    terrno = TAOS_SYSTEM_ERROR(errno);
Y
yihaoDeng 已提交
363 364
    taosCleanUpTcpClient(pClientObj);
    pClientObj = NULL;
F
freemine 已提交
365 366
  }
  return pClientObj;
J
Jeff Tao 已提交
367 368
}

陶建辉(Jeff)'s avatar
陶建辉(Jeff) 已提交
369
void taosStopTcpClient(void *chandle) {
Y
TD-2884  
yihaoDeng 已提交
370
  SClientObj *pClientObj = chandle;
陶建辉(Jeff)'s avatar
陶建辉(Jeff) 已提交
371

Y
TD-2884  
yihaoDeng 已提交
372 373
  if (pClientObj == NULL) return;

dengyihao's avatar
dengyihao 已提交
374
  tDebug("%s TCP client is stopped", pClientObj->label);
陶建辉(Jeff)'s avatar
陶建辉(Jeff) 已提交
375 376
}

J
Jeff Tao 已提交
377
void taosCleanUpTcpClient(void *chandle) {
Y
yihaoDeng 已提交
378 379
  SClientObj *pClientObj = chandle;
  if (pClientObj == NULL) return;
F
freemine 已提交
380
  for (int i = 0; i < pClientObj->numOfThreads; ++i) {
dengyihao's avatar
dengyihao 已提交
381
    SThreadObj *pThreadObj = pClientObj->pThreadObj[i];
F
freemine 已提交
382
    taosStopTcpThread(pThreadObj);
Y
yihaoDeng 已提交
383
  }
F
freemine 已提交
384

Y
yihaoDeng 已提交
385 386
  tDebug("%s TCP client is cleaned up", pClientObj->label);
  tfree(pClientObj->pThreadObj);
F
freemine 已提交
387
  tfree(pClientObj);
J
Jeff Tao 已提交
388 389
}

J
jtao1735 已提交
390
void *taosOpenTcpClientConnection(void *shandle, void *thandle, uint32_t ip, uint16_t port) {
dengyihao's avatar
dengyihao 已提交
391 392 393
  SClientObj *pClientObj = shandle;
  int32_t     index = atomic_load_32(&pClientObj->index) % pClientObj->numOfThreads;
  atomic_store_32(&pClientObj->index, index + 1);
Y
yihaoDeng 已提交
394
  SThreadObj *pThreadObj = pClientObj->pThreadObj[index];
J
Jeff Tao 已提交
395

S
TD-1057  
Shengliang Guan 已提交
396
  SOCKET fd = taosOpenTcpClientSocket(ip, port, pThreadObj->ip);
W
fix bug  
wpan 已提交
397 398 399
#if defined(_TD_WINDOWS_64) || defined(_TD_WINDOWS_32)
  if (fd == (SOCKET)-1) return NULL;
#else
400
  if (fd <= 0) return NULL;
W
fix bug  
wpan 已提交
401
#endif
J
Jeff Tao 已提交
402

陶建辉(Jeff)'s avatar
陶建辉(Jeff) 已提交
403
  struct sockaddr_in sin;
dengyihao's avatar
dengyihao 已提交
404 405 406
  uint16_t           localPort = 0;
  unsigned int       addrlen = sizeof(sin);
  if (getsockname(fd, (struct sockaddr *)&sin, &addrlen) == 0 && sin.sin_family == AF_INET && addrlen == sizeof(sin)) {
陶建辉(Jeff)'s avatar
陶建辉(Jeff) 已提交
407 408 409
    localPort = (uint16_t)ntohs(sin.sin_port);
  }

J
Jeff Tao 已提交
410
  SFdObj *pFdObj = taosMallocFdObj(pThreadObj, fd);
F
freemine 已提交
411

J
Jeff Tao 已提交
412 413 414
  if (pFdObj) {
    pFdObj->thandle = thandle;
    pFdObj->port = port;
J
jtao1735 已提交
415
    pFdObj->ip = ip;
H
Haojun Liao 已提交
416 417 418 419 420

    char ipport[40] = {0};
    taosIpPort2String(ip, port, ipport);
    tDebug("%s %p TCP connection to %s is created, localPort:%hu FD:%p numOfFds:%d", pThreadObj->label, thandle,
           ipport, localPort, pFdObj, pThreadObj->numOfFds);
J
Jeff Tao 已提交
421 422
  } else {
    tError("%s failed to malloc client FdObj(%s)", pThreadObj->label, strerror(errno));
H
Haojun Liao 已提交
423
    taosCloseSocket(fd);
J
Jeff Tao 已提交
424 425 426 427 428 429
  }

  return pFdObj;
}

void taosCloseTcpConnection(void *chandle) {
430
  SFdObj *pFdObj = chandle;
陶建辉(Jeff)'s avatar
TD-1645  
陶建辉(Jeff) 已提交
431
  if (pFdObj == NULL || pFdObj->signature != pFdObj) return;
432

陶建辉(Jeff)'s avatar
陶建辉(Jeff) 已提交
433
  SThreadObj *pThreadObj = pFdObj->pThreadObj;
F
freemine 已提交
434
  tDebug("%s %p TCP connection will be closed, FD:%p", pThreadObj->label, pFdObj->thandle, pFdObj);
陶建辉(Jeff)'s avatar
陶建辉(Jeff) 已提交
435 436

  // pFdObj->thandle = NULL;
陶建辉(Jeff)'s avatar
陶建辉(Jeff) 已提交
437
  pFdObj->closedByApp = 1;
S
Shengliang Guan 已提交
438
  taosShutDownSocketWR(pFdObj->fd);
439 440
}

J
Jeff Tao 已提交
441
int taosSendTcpData(uint32_t ip, uint16_t port, void *data, int len, void *chandle) {
442
  SFdObj *pFdObj = chandle;
陶建辉(Jeff)'s avatar
TD-1645  
陶建辉(Jeff) 已提交
443
  if (pFdObj == NULL || pFdObj->signature != pFdObj) return -1;
陶建辉(Jeff)'s avatar
陶建辉(Jeff) 已提交
444 445 446
  SThreadObj *pThreadObj = pFdObj->pThreadObj;

  int ret = taosWriteMsg(pFdObj->fd, data, len);
F
freemine 已提交
447
  tTrace("%s %p TCP data is sent, FD:%p fd:%d bytes:%d", pThreadObj->label, pFdObj->thandle, pFdObj, pFdObj->fd, ret);
448

陶建辉(Jeff)'s avatar
陶建辉(Jeff) 已提交
449
  return ret;
450 451
}

452 453 454 455
static void taosReportBrokenLink(SFdObj *pFdObj) {
  SThreadObj *pThreadObj = pFdObj->pThreadObj;

  // notify the upper layer, so it will clean the associated context
陶建辉(Jeff)'s avatar
陶建辉(Jeff) 已提交
456
  if (pFdObj->closedByApp == 0) {
S
Shengliang Guan 已提交
457
    taosShutDownSocketWR(pFdObj->fd);
陶建辉(Jeff)'s avatar
陶建辉(Jeff) 已提交
458

459 460 461 462 463 464
    SRecvInfo recvInfo;
    recvInfo.msg = NULL;
    recvInfo.msgLen = 0;
    recvInfo.ip = 0;
    recvInfo.port = 0;
    recvInfo.shandle = pThreadObj->shandle;
465
    recvInfo.thandle = pFdObj->thandle;
466 467 468
    recvInfo.chandle = NULL;
    recvInfo.connType = RPC_CONN_TCP;
    (*(pThreadObj->processData))(&recvInfo);
F
freemine 已提交
469
  }
陶建辉(Jeff)'s avatar
陶建辉(Jeff) 已提交
470 471 472 473 474

  taosFreeFdObj(pFdObj);
}

static int taosReadTcpData(SFdObj *pFdObj, SRecvInfo *pInfo) {
dengyihao's avatar
dengyihao 已提交
475 476 477
  SRpcHead rpcHead;
  int32_t  msgLen, leftLen, retLen, headLen;
  char *   buffer, *msg;
陶建辉(Jeff)'s avatar
陶建辉(Jeff) 已提交
478 479 480 481 482

  SThreadObj *pThreadObj = pFdObj->pThreadObj;

  headLen = taosReadMsg(pFdObj->fd, &rpcHead, sizeof(SRpcHead));
  if (headLen != sizeof(SRpcHead)) {
陶建辉(Jeff)'s avatar
陶建辉(Jeff) 已提交
483
    tDebug("%s %p read error, FD:%p headLen:%d", pThreadObj->label, pFdObj->thandle, pFdObj, headLen);
F
freemine 已提交
484
    return -1;
陶建辉(Jeff)'s avatar
陶建辉(Jeff) 已提交
485 486 487
  }

  msgLen = (int32_t)htonl((uint32_t)rpcHead.msgLen);
S
TD-1762  
Shengliang Guan 已提交
488 489
  int32_t size = msgLen + tsRpcOverhead;
  buffer = malloc(size);
S
Shengliang Guan 已提交
490
  if (NULL == buffer) {
陶建辉(Jeff)'s avatar
陶建辉(Jeff) 已提交
491
    tError("%s %p TCP malloc(size:%d) fail", pThreadObj->label, pFdObj->thandle, msgLen);
陶建辉(Jeff)'s avatar
陶建辉(Jeff) 已提交
492
    return -1;
陶建辉(Jeff)'s avatar
陶建辉(Jeff) 已提交
493
  } else {
dengyihao's avatar
dengyihao 已提交
494 495
    tTrace("%s %p read data, FD:%p fd:%d TCP malloc mem:%p", pThreadObj->label, pFdObj->thandle, pFdObj, pFdObj->fd,
           buffer);
陶建辉(Jeff)'s avatar
陶建辉(Jeff) 已提交
496 497 498 499 500 501 502
  }

  msg = buffer + tsRpcOverhead;
  leftLen = msgLen - headLen;
  retLen = taosReadMsg(pFdObj->fd, msg + headLen, leftLen);

  if (leftLen != retLen) {
dengyihao's avatar
dengyihao 已提交
503
    tError("%s %p read error, leftLen:%d retLen:%d FD:%p", pThreadObj->label, pFdObj->thandle, leftLen, retLen, pFdObj);
陶建辉(Jeff)'s avatar
陶建辉(Jeff) 已提交
504 505
    free(buffer);
    return -1;
506
  }
陶建辉(Jeff)'s avatar
陶建辉(Jeff) 已提交
507 508

  memcpy(msg, &rpcHead, sizeof(SRpcHead));
F
freemine 已提交
509

陶建辉(Jeff)'s avatar
陶建辉(Jeff) 已提交
510 511 512 513 514
  pInfo->msg = msg;
  pInfo->msgLen = msgLen;
  pInfo->ip = pFdObj->ip;
  pInfo->port = pFdObj->port;
  pInfo->shandle = pThreadObj->shandle;
515
  pInfo->thandle = pFdObj->thandle;
陶建辉(Jeff)'s avatar
陶建辉(Jeff) 已提交
516 517 518 519
  pInfo->chandle = pFdObj;
  pInfo->connType = RPC_CONN_TCP;

  if (pFdObj->closedByApp) {
F
freemine 已提交
520
    free(buffer);
陶建辉(Jeff)'s avatar
陶建辉(Jeff) 已提交
521 522 523 524
    return -1;
  }

  return 0;
525 526
}

J
Jeff Tao 已提交
527 528 529
#define maxEvents 10

static void *taosProcessTcpData(void *param) {
dengyihao's avatar
dengyihao 已提交
530 531
  SThreadObj *       pThreadObj = param;
  SFdObj *           pFdObj;
H
hzcheng 已提交
532
  struct epoll_event events[maxEvents];
533
  SRecvInfo          recvInfo;
534

H
Haojun Liao 已提交
535 536
  char name[16] = {0};
  snprintf(name, tListLen(name), "%s-tcp", pThreadObj->label);
537
  setThreadName(name);
F
freemine 已提交
538

H
hzcheng 已提交
539
  while (1) {
S
TD-1057  
Shengliang Guan 已提交
540
    int fdNum = epoll_wait(pThreadObj->pollFd, events, maxEvents, TAOS_EPOLL_WAIT_TIME);
541
    if (pThreadObj->stop) {
542
      tDebug("%s TCP thread get stop event, exiting...", pThreadObj->label);
543 544
      break;
    }
H
hzcheng 已提交
545 546
    if (fdNum < 0) continue;

J
Jeff Tao 已提交
547
    for (int i = 0; i < fdNum; ++i) {
H
hzcheng 已提交
548 549 550
      pFdObj = events[i].data.ptr;

      if (events[i].events & EPOLLERR) {
551
        tDebug("%s %p FD:%p epoll errors", pThreadObj->label, pFdObj->thandle, pFdObj);
552
        taosReportBrokenLink(pFdObj);
H
hzcheng 已提交
553 554 555
        continue;
      }

陶建辉(Jeff)'s avatar
陶建辉(Jeff) 已提交
556
      if (events[i].events & EPOLLRDHUP) {
557
        tDebug("%s %p FD:%p RD hang up", pThreadObj->label, pFdObj->thandle, pFdObj);
558
        taosReportBrokenLink(pFdObj);
H
hzcheng 已提交
559 560 561
        continue;
      }

陶建辉(Jeff)'s avatar
陶建辉(Jeff) 已提交
562
      if (events[i].events & EPOLLHUP) {
563
        tDebug("%s %p FD:%p hang up", pThreadObj->label, pFdObj->thandle, pFdObj);
564
        taosReportBrokenLink(pFdObj);
565 566
        continue;
      }
H
hzcheng 已提交
567

陶建辉(Jeff)'s avatar
陶建辉(Jeff) 已提交
568
      if (taosReadTcpData(pFdObj, &recvInfo) < 0) {
F
freemine 已提交
569
        shutdown(pFdObj->fd, SHUT_WR);
H
hzcheng 已提交
570 571 572
        continue;
      }

573
      pFdObj->thandle = (*(pThreadObj->processData))(&recvInfo);
574
      if (pFdObj->thandle == NULL) taosFreeFdObj(pFdObj);
H
hzcheng 已提交
575
    }
陶建辉(Jeff)'s avatar
TD-1645  
陶建辉(Jeff) 已提交
576

F
freemine 已提交
577
    if (pThreadObj->stop) break;
陶建辉(Jeff)'s avatar
TD-1645  
陶建辉(Jeff) 已提交
578 579
  }

dengyihao's avatar
dengyihao 已提交
580
  if (pThreadObj->pollFd >= 0) {
S
TD-2837  
Shengliang Guan 已提交
581
    EpollClose(pThreadObj->pollFd);
582 583
    pThreadObj->pollFd = -1;
  }
陶建辉(Jeff)'s avatar
TD-1645  
陶建辉(Jeff) 已提交
584 585

  while (pThreadObj->pHead) {
586
    pFdObj = pThreadObj->pHead;
陶建辉(Jeff)'s avatar
TD-1645  
陶建辉(Jeff) 已提交
587
    pThreadObj->pHead = pFdObj->next;
陶建辉(Jeff)'s avatar
TD-1669  
陶建辉(Jeff) 已提交
588
    taosReportBrokenLink(pFdObj);
H
hzcheng 已提交
589
  }
J
Jeff Tao 已提交
590

陶建辉(Jeff)'s avatar
TD-1645  
陶建辉(Jeff) 已提交
591
  pthread_mutex_destroy(&(pThreadObj->mutex));
陶建辉(Jeff)'s avatar
TD-1645  
陶建辉(Jeff) 已提交
592
  tDebug("%s TCP thread exits ...", pThreadObj->label);
S
TD-1848  
Shengliang Guan 已提交
593
  tfree(pThreadObj);
陶建辉(Jeff)'s avatar
TD-1645  
陶建辉(Jeff) 已提交
594

J
Jeff Tao 已提交
595
  return NULL;
H
hzcheng 已提交
596 597
}

S
TD-1057  
Shengliang Guan 已提交
598
static SFdObj *taosMallocFdObj(SThreadObj *pThreadObj, SOCKET fd) {
H
hzcheng 已提交
599 600
  struct epoll_event event;

J
Jeff Tao 已提交
601
  SFdObj *pFdObj = (SFdObj *)calloc(sizeof(SFdObj), 1);
陶建辉(Jeff)'s avatar
陶建辉(Jeff) 已提交
602 603 604
  if (pFdObj == NULL) {
    return NULL;
  }
H
hzcheng 已提交
605

陶建辉(Jeff)'s avatar
陶建辉(Jeff) 已提交
606
  pFdObj->closedByApp = 0;
J
Jeff Tao 已提交
607 608 609
  pFdObj->fd = fd;
  pFdObj->pThreadObj = pThreadObj;
  pFdObj->signature = pFdObj;
H
hzcheng 已提交
610

陶建辉(Jeff)'s avatar
陶建辉(Jeff) 已提交
611
  event.events = EPOLLIN | EPOLLRDHUP;
J
Jeff Tao 已提交
612 613
  event.data.ptr = pFdObj;
  if (epoll_ctl(pThreadObj->pollFd, EPOLL_CTL_ADD, fd, &event) < 0) {
S
TD-1848  
Shengliang Guan 已提交
614
    tfree(pFdObj);
F
freemine 已提交
615
    terrno = TAOS_SYSTEM_ERROR(errno);
J
Jeff Tao 已提交
616
    return NULL;
H
hzcheng 已提交
617 618
  }

J
Jeff Tao 已提交
619 620 621 622 623 624 625
  // notify the data process, add into the FdObj list
  pthread_mutex_lock(&(pThreadObj->mutex));
  pFdObj->next = pThreadObj->pHead;
  if (pThreadObj->pHead) (pThreadObj->pHead)->prev = pFdObj;
  pThreadObj->pHead = pFdObj;
  pThreadObj->numOfFds++;
  pthread_mutex_unlock(&(pThreadObj->mutex));
H
hzcheng 已提交
626

J
Jeff Tao 已提交
627
  return pFdObj;
H
hzcheng 已提交
628 629
}

J
Jeff Tao 已提交
630
static void taosFreeFdObj(SFdObj *pFdObj) {
631 632 633
  if (pFdObj == NULL) return;
  if (pFdObj->signature != pFdObj) return;

634
  SThreadObj *pThreadObj = pFdObj->pThreadObj;
J
Jeff Tao 已提交
635
  pthread_mutex_lock(&pThreadObj->mutex);
636

J
Jeff Tao 已提交
637 638 639 640 641 642
  if (pFdObj->signature == NULL) {
    pthread_mutex_unlock(&pThreadObj->mutex);
    return;
  }

  pFdObj->signature = NULL;
643
  epoll_ctl(pThreadObj->pollFd, EPOLL_CTL_DEL, pFdObj->fd, NULL);
陶建辉(Jeff)'s avatar
陶建辉(Jeff) 已提交
644
  taosCloseSocket(pFdObj->fd);
645 646 647

  pThreadObj->numOfFds--;
  if (pThreadObj->numOfFds < 0)
dengyihao's avatar
dengyihao 已提交
648 649
    tError("%s %p TCP thread:%d, number of FDs is negative!!!", pThreadObj->label, pFdObj->thandle,
           pThreadObj->threadId);
650 651 652 653 654 655 656 657 658 659 660

  if (pFdObj->prev) {
    (pFdObj->prev)->next = pFdObj->next;
  } else {
    pThreadObj->pHead = pFdObj->next;
  }

  if (pFdObj->next) {
    (pFdObj->next)->prev = pFdObj->prev;
  }

J
Jeff Tao 已提交
661
  pthread_mutex_unlock(&pThreadObj->mutex);
662

dengyihao's avatar
dengyihao 已提交
663 664
  tDebug("%s %p TCP connection is closed, FD:%p fd:%d numOfFds:%d", pThreadObj->label, pFdObj->thandle, pFdObj,
         pFdObj->fd, pThreadObj->numOfFds);
665

S
TD-1848  
Shengliang Guan 已提交
666
  tfree(pFdObj);
667
}
dengyihao's avatar
dengyihao 已提交
668 669

#endif