udfd.c 44.3 KB
Newer Older
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/>.
 */
dengyihao's avatar
dengyihao 已提交
15 16

// clang-format off
17 18
#include "uv.h"
#include "os.h"
S
slzhou 已提交
19 20
#include "fnLog.h"
#include "thash.h"
S
shenglian zhou 已提交
21

22 23
#include "tudf.h"
#include "tudfInt.h"
D
dapan1121 已提交
24
#include "version.h"
25

26
#include "tdatablock.h"
27 28 29 30
#include "tdataformat.h"
#include "tglobal.h"
#include "tmsg.h"
#include "trpc.h"
31
#include "tmisce.h"
32
// clang-format on
33

S
slzhou 已提交
34 35 36
#define UDFD_MAX_SCRIPT_PLUGINS 64
#define UDFD_MAX_SCRIPT_TYPE    1
#define UDFD_MAX_PLUGIN_FUNCS   9
37 38 39 40 41 42 43 44 45 46 47 48 49 50 51

typedef struct SUdfCPluginCtx {
  uv_lib_t lib;

  TUdfScalarProcFunc scalarProcFunc;

  TUdfAggStartFunc   aggStartFunc;
  TUdfAggProcessFunc aggProcFunc;
  TUdfAggFinishFunc  aggFinishFunc;
  TUdfAggMergeFunc   aggMergeFunc;

  TUdfInitFunc    initFunc;
  TUdfDestroyFunc destroyFunc;
} SUdfCPluginCtx;

52
int32_t udfdCPluginOpen(SScriptUdfEnvItem *items, int numItems) { return 0; }
53

54
int32_t udfdCPluginClose() { return 0; }
55

S
slzhou 已提交
56
int32_t udfdCPluginUdfInit(SScriptUdfInfo *udf, void **pUdfCtx) {
57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76
  int32_t         err = 0;
  SUdfCPluginCtx *udfCtx = taosMemoryCalloc(1, sizeof(SUdfCPluginCtx));
  err = uv_dlopen(udf->path, &udfCtx->lib);
  if (err != 0) {
    fnError("can not load library %s. error: %s", udf->path, uv_strerror(err));
    return TSDB_CODE_UDF_LOAD_UDF_FAILURE;
  }
  const char *udfName = udf->name;
  char        initFuncName[TSDB_FUNC_NAME_LEN + 5] = {0};
  char       *initSuffix = "_init";
  strcpy(initFuncName, udfName);
  strncat(initFuncName, initSuffix, strlen(initSuffix));
  uv_dlsym(&udfCtx->lib, initFuncName, (void **)(&udfCtx->initFunc));

  char  destroyFuncName[TSDB_FUNC_NAME_LEN + 5] = {0};
  char *destroySuffix = "_destroy";
  strcpy(destroyFuncName, udfName);
  strncat(destroyFuncName, destroySuffix, strlen(destroySuffix));
  uv_dlsym(&udfCtx->lib, destroyFuncName, (void **)(&udfCtx->destroyFunc));

S
slzhou 已提交
77
  if (udf->funcType == UDF_FUNC_TYPE_SCALAR) {
78 79 80
    char processFuncName[TSDB_FUNC_NAME_LEN] = {0};
    strcpy(processFuncName, udfName);
    uv_dlsym(&udfCtx->lib, processFuncName, (void **)(&udfCtx->scalarProcFunc));
S
slzhou 已提交
81
  } else if (udf->funcType == UDF_FUNC_TYPE_AGG) {
82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100
    char processFuncName[TSDB_FUNC_NAME_LEN] = {0};
    strcpy(processFuncName, udfName);
    uv_dlsym(&udfCtx->lib, processFuncName, (void **)(&udfCtx->aggProcFunc));
    char  startFuncName[TSDB_FUNC_NAME_LEN + 6] = {0};
    char *startSuffix = "_start";
    strncpy(startFuncName, processFuncName, sizeof(startFuncName));
    strncat(startFuncName, startSuffix, strlen(startSuffix));
    uv_dlsym(&udfCtx->lib, startFuncName, (void **)(&udfCtx->aggStartFunc));
    char  finishFuncName[TSDB_FUNC_NAME_LEN + 7] = {0};
    char *finishSuffix = "_finish";
    strncpy(finishFuncName, processFuncName, sizeof(finishFuncName));
    strncat(finishFuncName, finishSuffix, strlen(finishSuffix));
    uv_dlsym(&udfCtx->lib, finishFuncName, (void **)(&udfCtx->aggFinishFunc));
    char  mergeFuncName[TSDB_FUNC_NAME_LEN + 6] = {0};
    char *mergeSuffix = "_merge";
    strncpy(mergeFuncName, processFuncName, sizeof(mergeFuncName));
    strncat(mergeFuncName, mergeSuffix, strlen(mergeSuffix));
    uv_dlsym(&udfCtx->lib, mergeFuncName, (void **)(&udfCtx->aggMergeFunc));
  }
S
slzhou 已提交
101
  int32_t code = 0;
102
  if (udfCtx->initFunc) {
S
slzhou 已提交
103 104 105 106 107 108
    code = (udfCtx->initFunc)();
    if (code != 0) {
      uv_dlclose(&udfCtx->lib);
      taosMemoryFree(udfCtx);
      return code;
    }
109 110 111 112 113 114 115
  }
  *pUdfCtx = udfCtx;
  return 0;
}

int32_t udfdCPluginUdfDestroy(void *udfCtx) {
  SUdfCPluginCtx *ctx = udfCtx;
S
slzhou 已提交
116
  int32_t         code = 0;
117
  if (ctx->destroyFunc) {
S
slzhou 已提交
118
    code = (ctx->destroyFunc)();
119 120 121
  }
  uv_dlclose(&ctx->lib);
  taosMemoryFree(ctx);
S
slzhou 已提交
122
  return code;
123 124 125 126 127
}

int32_t udfdCPluginUdfScalarProc(SUdfDataBlock *block, SUdfColumn *resultCol, void *udfCtx) {
  SUdfCPluginCtx *ctx = udfCtx;
  if (ctx->scalarProcFunc) {
S
slzhou 已提交
128 129
    return ctx->scalarProcFunc(block, resultCol);
  } else {
S
slzhou 已提交
130 131
    fnError("udfd c plugin scalar proc not implemented");
    return TSDB_CODE_UDF_FUNC_EXEC_FAILURE;
132 133 134 135 136 137
  }
}

int32_t udfdCPluginUdfAggStart(SUdfInterBuf *buf, void *udfCtx) {
  SUdfCPluginCtx *ctx = udfCtx;
  if (ctx->aggStartFunc) {
S
slzhou 已提交
138 139
    return ctx->aggStartFunc(buf);
  } else {
S
slzhou 已提交
140 141
    fnError("udfd c plugin aggregation start not implemented");
    return TSDB_CODE_UDF_FUNC_EXEC_FAILURE;
142 143 144 145 146 147 148
  }
  return 0;
}

int32_t udfdCPluginUdfAggProc(SUdfDataBlock *block, SUdfInterBuf *interBuf, SUdfInterBuf *newInterBuf, void *udfCtx) {
  SUdfCPluginCtx *ctx = udfCtx;
  if (ctx->aggProcFunc) {
S
slzhou 已提交
149 150
    return ctx->aggProcFunc(block, interBuf, newInterBuf);
  } else {
S
slzhou 已提交
151 152
    fnError("udfd c plugin aggregation process not implemented");
    return TSDB_CODE_UDF_FUNC_EXEC_FAILURE;
153 154 155 156 157 158 159
  }
}

int32_t udfdCPluginUdfAggMerge(SUdfInterBuf *inputBuf1, SUdfInterBuf *inputBuf2, SUdfInterBuf *outputBuf,
                               void *udfCtx) {
  SUdfCPluginCtx *ctx = udfCtx;
  if (ctx->aggMergeFunc) {
S
slzhou 已提交
160 161
    return ctx->aggMergeFunc(inputBuf1, inputBuf2, outputBuf);
  } else {
S
slzhou 已提交
162 163
    fnError("udfd c plugin aggregation merge not implemented");
    return TSDB_CODE_UDF_FUNC_EXEC_FAILURE;
164 165 166 167 168 169
  }
}

int32_t udfdCPluginUdfAggFinish(SUdfInterBuf *buf, SUdfInterBuf *resultData, void *udfCtx) {
  SUdfCPluginCtx *ctx = udfCtx;
  if (ctx->aggFinishFunc) {
S
slzhou 已提交
170 171
    return ctx->aggFinishFunc(buf, resultData);
  } else {
S
slzhou 已提交
172 173
    fnError("udfd c plugin aggregation finish not implemented");
    return TSDB_CODE_UDF_FUNC_EXEC_FAILURE;
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
  }
  return 0;
}

// for c, the function pointer are filled directly and libloaded = true;
// for others, dlopen/dlsym to find function pointers
typedef struct SUdfScriptPlugin {
  int8_t scriptType;

  char     libPath[PATH_MAX];
  bool     libLoaded;
  uv_lib_t lib;

  TScriptUdfScalarProcFunc udfScalarProcFunc;
  TScriptUdfAggStartFunc   udfAggStartFunc;
  TScriptUdfAggProcessFunc udfAggProcFunc;
  TScriptUdfAggMergeFunc   udfAggMergeFunc;
  TScriptUdfAggFinishFunc  udfAggFinishFunc;

  TScriptUdfInitFunc    udfInitFunc;
  TScriptUdfDestoryFunc udfDestroyFunc;

  TScriptOpenFunc  openFunc;
  TScriptCloseFunc closeFunc;
} SUdfScriptPlugin;

S
slzhou 已提交
200
typedef struct SUdfdContext {
201
  uv_loop_t  *loop;
S
slzhou 已提交
202
  uv_pipe_t   ctrlPipe;
203
  uv_signal_t intrSignal;
204
  char        listenPipeName[PATH_MAX + UDF_LISTEN_PIPE_NAME_LEN + 2];
S
slzhou 已提交
205
  uv_pipe_t   listeningPipe;
S
slzhou 已提交
206

207 208 209
  void     *clientRpc;
  SCorEpSet mgmtEp;

S
slzhou 已提交
210
  uv_mutex_t udfsMutex;
211
  SHashObj  *udfsHash;
S
slzhou 已提交
212

213
  uv_mutex_t        scriptPluginsMutex;
S
slzhou 已提交
214
  SUdfScriptPlugin *scriptPlugins[UDFD_MAX_SCRIPT_PLUGINS];
215

216
  SArray *residentFuncs;
217

S
slzhou 已提交
218 219 220 221
  bool printVersion;
} SUdfdContext;

SUdfdContext global;
222

223 224 225
struct SUdfdUvConn;
struct SUvUdfWork;

226
typedef struct SUdfdUvConn {
S
slzhou 已提交
227
  uv_stream_t *client;
228
  char        *inputBuf;
S
slzhou 已提交
229 230 231
  int32_t      inputLen;
  int32_t      inputCap;
  int32_t      inputTotal;
232 233

  struct SUvUdfWork *pWorkList;  // head of work list
234 235 236
} SUdfdUvConn;

typedef struct SUvUdfWork {
237
  SUdfdUvConn *conn;
S
slzhou 已提交
238 239
  uv_buf_t     input;
  uv_buf_t     output;
240 241

  struct SUvUdfWork *pWorkNext;
242 243
} SUvUdfWork;

244
typedef enum { UDF_STATE_INIT = 0, UDF_STATE_LOADING, UDF_STATE_READY } EUdfState;
245

S
slzhou 已提交
246
typedef struct SUdf {
247
  char    name[TSDB_FUNC_NAME_LEN + 1];
248
  int32_t version;
S
slzhou 已提交
249

dengyihao's avatar
dengyihao 已提交
250 251 252
  int8_t  funcType;
  int8_t  scriptType;
  int8_t  outputType;
S
slzhou 已提交
253 254 255
  int32_t outputLen;
  int32_t bufSize;

dengyihao's avatar
dengyihao 已提交
256
  char path[PATH_MAX];
S
slzhou 已提交
257

258 259 260 261 262
  int32_t    refCount;
  EUdfState  state;
  uv_mutex_t lock;
  uv_cond_t  condReady;
  bool       resident;
263

264 265
  SUdfScriptPlugin *scriptPlugin;
  void             *scriptUdfCtx;
266 267 268

  int64_t lastFetchTime;  // last fetch time in milliseconds
  bool    expired;
269 270
} SUdf;

271
typedef struct SUdfcFuncHandle {
S
slzhou 已提交
272
  SUdf *udf;
273
} SUdfcFuncHandle;
274

275 276 277 278 279 280 281
typedef enum EUdfdRpcReqRspType {
  UDFD_RPC_MNODE_CONNECT = 0,
  UDFD_RPC_RETRIVE_FUNC,
} EUdfdRpcReqRspType;

typedef struct SUdfdRpcSendRecvInfo {
  EUdfdRpcReqRspType rpcType;
dengyihao's avatar
dengyihao 已提交
282
  int32_t            code;
283
  void              *param;
dengyihao's avatar
dengyihao 已提交
284
  uv_sem_t           resultSem;
285 286
} SUdfdRpcSendRecvInfo;

dengyihao's avatar
dengyihao 已提交
287
static void    udfdProcessRpcRsp(void *parent, SRpcMsg *pMsg, SEpSet *pEpSet);
S
shenglian zhou 已提交
288 289
static int32_t udfdFillUdfInfoFromMNode(void *clientRpc, char *udfName, SUdf *udf);
static int32_t udfdConnectToMnode();
dengyihao's avatar
dengyihao 已提交
290
static bool    udfdRpcRfp(int32_t code, tmsg_t msgType);
dengyihao's avatar
dengyihao 已提交
291
static int     initEpSetFromCfg(const char *firstEp, const char *secondEp, SCorEpSet *pEpSet);
S
shenglian zhou 已提交
292 293 294 295 296 297 298 299 300 301 302 303 304 305 306 307 308
static int32_t udfdOpenClientRpc();
static int32_t udfdCloseClientRpc();

static void udfdProcessSetupRequest(SUvUdfWork *uvUdf, SUdfRequest *request);
static void udfdProcessCallRequest(SUvUdfWork *uvUdf, SUdfRequest *request);
static void udfdProcessTeardownRequest(SUvUdfWork *uvUdf, SUdfRequest *request);
static void udfdProcessRequest(uv_work_t *req);
static void udfdOnWrite(uv_write_t *req, int status);
static void udfdSendResponse(uv_work_t *work, int status);
static void udfdAllocBuffer(uv_handle_t *handle, size_t suggestedSize, uv_buf_t *buf);
static bool isUdfdUvMsgComplete(SUdfdUvConn *pipe);
static void udfdHandleRequest(SUdfdUvConn *conn);
static void udfdPipeCloseCb(uv_handle_t *pipe);
static void udfdUvHandleError(SUdfdUvConn *conn) { uv_close((uv_handle_t *)conn->client, udfdPipeCloseCb); }
static void udfdPipeRead(uv_stream_t *client, ssize_t nread, const uv_buf_t *buf);
static void udfdOnNewConnection(uv_stream_t *server, int status);

dengyihao's avatar
dengyihao 已提交
309
static void    udfdIntrSignalHandler(uv_signal_t *handle, int signum);
S
shenglian zhou 已提交
310 311
static int32_t removeListeningPipe();

dengyihao's avatar
dengyihao 已提交
312
static void    udfdPrintVersion();
S
shenglian zhou 已提交
313 314 315
static int32_t udfdParseArgs(int32_t argc, char *argv[]);
static int32_t udfdInitLog();

dengyihao's avatar
dengyihao 已提交
316 317
static void    udfdCtrlAllocBufCb(uv_handle_t *handle, size_t suggested_size, uv_buf_t *buf);
static void    udfdCtrlReadCb(uv_stream_t *q, ssize_t nread, const uv_buf_t *buf);
S
shenglian zhou 已提交
318
static int32_t udfdUvInit();
dengyihao's avatar
dengyihao 已提交
319
static void    udfdCloseWalkCb(uv_handle_t *handle, void *arg);
S
shenglian zhou 已提交
320
static int32_t udfdRun();
dengyihao's avatar
dengyihao 已提交
321
static void    udfdConnectMnodeThreadFunc(void *args);
S
shenglian zhou 已提交
322

323 324
SUdf *udfdNewUdf(const char *udfName);
void  udfdInitializeCPlugin(SUdfScriptPlugin *plugin) {
325 326 327 328 329 330 331 332 333 334
  plugin->scriptType = TSDB_FUNC_SCRIPT_BIN_LIB;
  plugin->openFunc = udfdCPluginOpen;
  plugin->closeFunc = udfdCPluginClose;
  plugin->udfInitFunc = udfdCPluginUdfInit;
  plugin->udfDestroyFunc = udfdCPluginUdfDestroy;
  plugin->udfScalarProcFunc = udfdCPluginUdfScalarProc;
  plugin->udfAggStartFunc = udfdCPluginUdfAggStart;
  plugin->udfAggProcFunc = udfdCPluginUdfAggProc;
  plugin->udfAggMergeFunc = udfdCPluginUdfAggMerge;
  plugin->udfAggFinishFunc = udfdCPluginUdfAggFinish;
335

S
slzhou 已提交
336 337
  SScriptUdfEnvItem items[1] = {{"LD_LIBRARY_PATH", tsUdfdLdLibPath}};
  plugin->openFunc(items, 1);
338 339 340 341 342 343 344 345 346 347 348 349 350 351 352 353 354 355 356
  return;
}

int32_t udfdLoadSharedLib(char *libPath, uv_lib_t *pLib, const char *funcName[], void **func[], int numOfFuncs) {
  int err = uv_dlopen(libPath, pLib);
  if (err != 0) {
    fnError("can not load library %s. error: %s", libPath, uv_strerror(err));
    return TSDB_CODE_UDF_LOAD_UDF_FAILURE;
  }

  for (int i = 0; i < numOfFuncs; ++i) {
    err = uv_dlsym(pLib, funcName[i], func[i]);
    if (err != 0) {
      fnError("load library function failed. lib %s function %s", libPath, funcName[i]);
    }
  }
  return 0;
}

357
int32_t udfdInitializePythonPlugin(SUdfScriptPlugin *plugin) {
358
  plugin->scriptType = TSDB_FUNC_SCRIPT_PYTHON;
S
slzhou 已提交
359
  // todo: windows support
S
slzhou 已提交
360
  sprintf(plugin->libPath, "%s", "libtaospyudf.so");
361
  plugin->libLoaded = false;
S
slzhou 已提交
362 363 364 365
  const char *funcName[UDFD_MAX_PLUGIN_FUNCS] = {"pyOpen",         "pyClose",         "pyUdfInit",
                                                 "pyUdfDestroy",   "pyUdfScalarProc", "pyUdfAggStart",
                                                 "pyUdfAggFinish", "pyUdfAggProc",    "pyUdfAggMerge"};
  void      **funcs[UDFD_MAX_PLUGIN_FUNCS] = {
366 367 368
      (void **)&plugin->openFunc,         (void **)&plugin->closeFunc,         (void **)&plugin->udfInitFunc,
      (void **)&plugin->udfDestroyFunc,   (void **)&plugin->udfScalarProcFunc, (void **)&plugin->udfAggStartFunc,
      (void **)&plugin->udfAggFinishFunc, (void **)&plugin->udfAggProcFunc,    (void **)&plugin->udfAggMergeFunc};
S
slzhou 已提交
369
  int32_t err = udfdLoadSharedLib(plugin->libPath, &plugin->lib, funcName, funcs, UDFD_MAX_PLUGIN_FUNCS);
370 371
  if (err != 0) {
    fnError("can not load python plugin. lib path %s", plugin->libPath);
372
    return err;
373
  }
374

375
  if (plugin->openFunc) {
376 377 378 379 380 381 382
    int16_t lenPythonPath = strlen(tsUdfdLdLibPath) + strlen(tsTempDir) + 1 + 1;  // tsTempDir:tsUdfdLdLibPath
    char   *pythonPath = taosMemoryMalloc(lenPythonPath);
#ifdef WINDOWS
    snprintf(pythonPath, lenPythonPath, "%s;%s", tsTempDir, tsUdfdLdLibPath);
#else
    snprintf(pythonPath, lenPythonPath, "%s:%s", tsTempDir, tsUdfdLdLibPath);
#endif
S
slzhou 已提交
383
    SScriptUdfEnvItem items[] = {{"PYTHONPATH", pythonPath}, {"LOGDIR", tsLogDir}};
S
slzhou 已提交
384
    err = plugin->openFunc(items, 2);
385
    taosMemoryFree(pythonPath);
386
  }
387 388 389 390 391
  if (err != 0) {
    fnError("udf script python plugin open func failed. error: %d", err);
    uv_dlclose(&plugin->lib);
    return err;
  }
392
  plugin->libLoaded = true;
393 394

  return 0;
395 396
}

397
void udfdDeinitCPlugin(SUdfScriptPlugin *plugin) {
398 399 400 401 402 403 404 405 406 407 408
  if (plugin->closeFunc) {
    plugin->closeFunc();
  }
  plugin->openFunc = NULL;
  plugin->closeFunc = NULL;
  plugin->udfInitFunc = NULL;
  plugin->udfDestroyFunc = NULL;
  plugin->udfScalarProcFunc = NULL;
  plugin->udfAggStartFunc = NULL;
  plugin->udfAggProcFunc = NULL;
  plugin->udfAggMergeFunc = NULL;
409
  plugin->udfAggFinishFunc = NULL;
410 411 412 413
  return;
}

void udfdDeinitPythonPlugin(SUdfScriptPlugin *plugin) {
414 415 416
  if (plugin->closeFunc) {
    plugin->closeFunc();
  }
S
slzhou 已提交
417
  uv_dlclose(&plugin->lib);
418 419 420
  if (plugin->libLoaded) {
    plugin->libLoaded = false;
  }
421 422 423 424 425 426 427 428 429
  plugin->openFunc = NULL;
  plugin->closeFunc = NULL;
  plugin->udfInitFunc = NULL;
  plugin->udfDestroyFunc = NULL;
  plugin->udfScalarProcFunc = NULL;
  plugin->udfAggStartFunc = NULL;
  plugin->udfAggProcFunc = NULL;
  plugin->udfAggMergeFunc = NULL;
  plugin->udfAggFinishFunc = NULL;
430 431
}

S
slzhou 已提交
432 433
int32_t udfdInitScriptPlugin(int8_t scriptType) {
  SUdfScriptPlugin *plugin = taosMemoryCalloc(1, sizeof(SUdfScriptPlugin));
434

S
slzhou 已提交
435 436 437 438
  switch (scriptType) {
    case TSDB_FUNC_SCRIPT_BIN_LIB:
      udfdInitializeCPlugin(plugin);
      break;
439 440 441 442 443 444
    case TSDB_FUNC_SCRIPT_PYTHON: {
      int32_t err = udfdInitializePythonPlugin(plugin);
      if (err != 0) {
        taosMemoryFree(plugin);
        return err;
      }
S
slzhou 已提交
445
      break;
446
    }
S
slzhou 已提交
447 448 449 450 451
    default:
      fnError("udf script type %d not supported", scriptType);
      taosMemoryFree(plugin);
      return TSDB_CODE_UDF_SCRIPT_NOT_SUPPORTED;
  }
452

S
slzhou 已提交
453 454
  global.scriptPlugins[scriptType] = plugin;
  return TSDB_CODE_SUCCESS;
455 456 457
}

void udfdDeinitScriptPlugins() {
458 459
  SUdfScriptPlugin *plugin = NULL;
  plugin = global.scriptPlugins[TSDB_FUNC_SCRIPT_PYTHON];
S
slzhou 已提交
460 461 462 463
  if (plugin != NULL) {
    udfdDeinitPythonPlugin(plugin);
    taosMemoryFree(plugin);
  }
464 465

  plugin = global.scriptPlugins[TSDB_FUNC_SCRIPT_BIN_LIB];
S
slzhou 已提交
466 467 468 469
  if (plugin != NULL) {
    udfdDeinitCPlugin(plugin);
    taosMemoryFree(plugin);
  }
470 471 472
  return;
}

S
shenglian zhou 已提交
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
void udfdProcessRequest(uv_work_t *req) {
  SUvUdfWork *uvUdf = (SUvUdfWork *)(req->data);
  SUdfRequest request = {0};
  decodeUdfRequest(uvUdf->input.base, &request);

  switch (request.type) {
    case UDF_TASK_SETUP: {
      udfdProcessSetupRequest(uvUdf, &request);
      break;
    }

    case UDF_TASK_CALL: {
      udfdProcessCallRequest(uvUdf, &request);
      break;
    }
    case UDF_TASK_TEARDOWN: {
      udfdProcessTeardownRequest(uvUdf, &request);
      break;
    }
    default: {
      break;
    }
  }
}

S
slzhou 已提交
498
void convertUdf2UdfInfo(SUdf *udf, SScriptUdfInfo *udfInfo) {
499
  udfInfo->bufSize = udf->bufSize;
S
slzhou 已提交
500 501 502 503 504
  if (udf->funcType == TSDB_FUNC_TYPE_AGGREGATE) {
    udfInfo->funcType = UDF_FUNC_TYPE_AGG;
  } else if (udf->funcType == TSDB_FUNC_TYPE_SCALAR) {
    udfInfo->funcType = UDF_FUNC_TYPE_SCALAR;
  }
S
slzhou 已提交
505
  udfInfo->name = udf->name;
506
  udfInfo->version = udf->version;
507 508
  udfInfo->outputLen = udf->outputLen;
  udfInfo->outputType = udf->outputType;
S
slzhou 已提交
509
  udfInfo->path = udf->path;
510 511 512
  udfInfo->scriptType = udf->scriptType;
}

S
slzhou 已提交
513 514 515
int32_t udfdRenameUdfFile(SUdf *udf) {
  char newPath[PATH_MAX];
  if (udf->scriptType == TSDB_FUNC_SCRIPT_BIN_LIB) {
S
slzhou 已提交
516
    snprintf(newPath, PATH_MAX, "%s/lib%s_%d_%" PRId64 ".so", tsTempDir, udf->name, udf->version, udf->lastFetchTime);
S
slzhou 已提交
517
  } else if (udf->scriptType == TSDB_FUNC_SCRIPT_PYTHON) {
S
slzhou 已提交
518
    snprintf(newPath, PATH_MAX, "%s/%s_%d_%" PRId64 ".py", tsTempDir, udf->name, udf->version, udf->lastFetchTime);
S
slzhou 已提交
519 520 521
  } else {
    return TSDB_CODE_UDF_SCRIPT_NOT_SUPPORTED;
  }
522 523 524 525
  int32_t code = taosRenameFile(udf->path, newPath);
  if (code == 0) {
    sprintf(udf->path, "%s", newPath);
  }
S
slzhou 已提交
526 527 528
  return 0;
}

529 530 531 532 533 534 535
int32_t udfdInitUdf(char *udfName, SUdf *udf) {
  int32_t err = 0;
  err = udfdFillUdfInfoFromMNode(global.clientRpc, udfName, udf);
  if (err != 0) {
    fnError("can not retrieve udf from mnode. udf name %s", udfName);
    return TSDB_CODE_UDF_LOAD_UDF_FAILURE;
  }
S
slzhou 已提交
536 537 538 539 540
  if (udf->scriptType > UDFD_MAX_SCRIPT_TYPE) {
    fnError("udf name %s script type %d not supported", udfName, udf->scriptType);
    return TSDB_CODE_UDF_SCRIPT_NOT_SUPPORTED;
  }

541 542
  uv_mutex_lock(&global.scriptPluginsMutex);
  SUdfScriptPlugin *scriptPlugin = global.scriptPlugins[udf->scriptType];
S
slzhou 已提交
543
  if (scriptPlugin == NULL) {
S
slzhou 已提交
544 545 546 547 548
    err = udfdInitScriptPlugin(udf->scriptType);
    if (err != 0) {
      uv_mutex_unlock(&global.scriptPluginsMutex);
      return err;
    }
549 550
  }
  uv_mutex_unlock(&global.scriptPluginsMutex);
S
slzhou 已提交
551 552 553 554
  udf->scriptPlugin = global.scriptPlugins[udf->scriptType];

  udfdRenameUdfFile(udf);

S
slzhou 已提交
555
  SScriptUdfInfo info = {0};
556
  convertUdf2UdfInfo(udf, &info);
S
slzhou 已提交
557 558 559 560 561
  err = udf->scriptPlugin->udfInitFunc(&info, &udf->scriptUdfCtx);
  if (err != 0) {
    fnError("udf name %s init failed. error %d", udfName, err);
    return err;
  }
562

563
  fnInfo("udf init succeeded. name %s type %d context %p", udf->name, udf->scriptType, (void *)udf->scriptUdfCtx);
564 565 566
  return 0;
}

567 568 569 570 571 572 573 574 575 576 577 578 579 580 581 582 583 584 585 586 587
SUdf *udfdNewUdf(const char *udfName) {
  SUdf *udfNew = taosMemoryCalloc(1, sizeof(SUdf));
  udfNew->refCount = 1;
  strncpy(udfNew->name, udfName, TSDB_FUNC_NAME_LEN);

  udfNew->state = UDF_STATE_INIT;
  uv_mutex_init(&udfNew->lock);
  uv_cond_init(&udfNew->condReady);

  udfNew->resident = false;
  udfNew->expired = false;
  for (int32_t i = 0; i < taosArrayGetSize(global.residentFuncs); ++i) {
    char *funcName = taosArrayGet(global.residentFuncs, i);
    if (strcmp(udfName, funcName) == 0) {
      udfNew->resident = true;
      break;
    }
  }
  return udfNew;
}

588
SUdf *udfdGetOrCreateUdf(const char *udfName) {
S
shenglian zhou 已提交
589
  uv_mutex_lock(&global.udfsMutex);
590 591
  SUdf  **pUdfHash = taosHashGet(global.udfsHash, udfName, strlen(udfName));
  int64_t currTime = taosGetTimestampSec();
S
slzhou 已提交
592 593 594 595 596 597 598 599 600 601 602 603
  bool    expired = false;
  if (pUdfHash) {
    expired = currTime - (*pUdfHash)->lastFetchTime > 10 * 1000;
    if (!expired) {
      ++(*pUdfHash)->refCount;
      SUdf *udf = *pUdfHash;
      uv_mutex_unlock(&global.udfsMutex);
      return udf;
    } else {
      (*pUdfHash)->expired = true;
      taosHashRemove(global.udfsHash, udfName, strlen(udfName));
    }
604 605 606 607 608 609 610
  }

  SUdf  *udf = udfdNewUdf(udfName);
  SUdf **pUdf = &udf;
  taosHashPut(global.udfsHash, udfName, strlen(udfName), pUdf, POINTER_BYTES);
  uv_mutex_unlock(&global.udfsMutex);

611 612 613 614 615 616 617 618 619
  return udf;
}

void udfdProcessSetupRequest(SUvUdfWork *uvUdf, SUdfRequest *request) {
  // TODO: tracable id from client. connect, setup, call, teardown
  fnInfo("setup request. seq num: %" PRId64 ", udf name: %s", request->seqNum, request->setup.udfName);

  SUdfSetupRequest *setup = &request->setup;
  int32_t           code = TSDB_CODE_SUCCESS;
620
  SUdf             *udf = NULL;
621 622

  udf = udfdGetOrCreateUdf(setup->udfName);
S
shenglian zhou 已提交
623 624 625 626

  uv_mutex_lock(&udf->lock);
  if (udf->state == UDF_STATE_INIT) {
    udf->state = UDF_STATE_LOADING;
627
    code = udfdInitUdf(setup->udfName, udf);
628 629 630 631
    if (code == 0) {
      udf->state = UDF_STATE_READY;
    } else {
      udf->state = UDF_STATE_INIT;
S
slzhou 已提交
632
    }
S
shenglian zhou 已提交
633 634 635
    uv_cond_broadcast(&udf->condReady);
    uv_mutex_unlock(&udf->lock);
  } else {
636
    while (udf->state == UDF_STATE_LOADING) {
S
shenglian zhou 已提交
637 638 639 640 641 642 643 644 645 646
      uv_cond_wait(&udf->condReady, &udf->lock);
    }
    uv_mutex_unlock(&udf->lock);
  }
  SUdfcFuncHandle *handle = taosMemoryMalloc(sizeof(SUdfcFuncHandle));
  handle->udf = udf;

  SUdfResponse rsp;
  rsp.seqNum = request->seqNum;
  rsp.type = request->type;
S
slzhou 已提交
647
  rsp.code = (code != 0) ? TSDB_CODE_UDF_FUNC_EXEC_FAILURE : 0;
S
shenglian zhou 已提交
648 649
  rsp.setupRsp.udfHandle = (int64_t)(handle);
  rsp.setupRsp.outputType = udf->outputType;
S
slzhou 已提交
650
  rsp.setupRsp.bytes = udf->outputLen;
S
shenglian zhou 已提交
651 652 653 654 655 656 657 658 659 660 661 662 663 664 665 666
  rsp.setupRsp.bufSize = udf->bufSize;

  int32_t len = encodeUdfResponse(NULL, &rsp);
  rsp.msgLen = len;
  void *bufBegin = taosMemoryMalloc(len);
  void *buf = bufBegin;
  encodeUdfResponse(&buf, &rsp);

  uvUdf->output = uv_buf_init(bufBegin, len);

  taosMemoryFree(uvUdf->input.base);
  return;
}

void udfdProcessCallRequest(SUvUdfWork *uvUdf, SUdfRequest *request) {
  SUdfCallRequest *call = &request->call;
667 668 669 670
  fnDebug("call request. call type %d, handle: %" PRIx64 ", seq num %" PRId64, call->callType, call->udfHandle,
          request->seqNum);
  SUdfcFuncHandle  *handle = (SUdfcFuncHandle *)(call->udfHandle);
  SUdf             *udf = handle->udf;
S
shenglian zhou 已提交
671
  SUdfResponse      response = {0};
672
  SUdfResponse     *rsp = &response;
S
shenglian zhou 已提交
673 674 675 676 677 678
  SUdfCallResponse *subRsp = &rsp->callRsp;

  int32_t code = TSDB_CODE_SUCCESS;
  switch (call->callType) {
    case TSDB_UDF_CALL_SCALA_PROC: {
      SUdfColumn output = {0};
S
slzhou 已提交
679
      output.colMeta.bytes = udf->outputLen;
680 681 682
      output.colMeta.type = udf->outputType;
      output.colMeta.precision = 0;
      output.colMeta.scale = 0;
S
shenglian zhou 已提交
683 684
      SUdfDataBlock input = {0};
      convertDataBlockToUdfDataBlock(&call->block, &input);
685
      code = udf->scriptPlugin->udfScalarProcFunc(&input, &output, udf->scriptUdfCtx);
S
shenglian zhou 已提交
686 687 688 689 690 691 692
      freeUdfDataDataBlock(&input);
      convertUdfColumnToDataBlock(&output, &response.callRsp.resultData);
      freeUdfColumn(&output);
      break;
    }
    case TSDB_UDF_CALL_AGG_INIT: {
      SUdfInterBuf outBuf = {.buf = taosMemoryMalloc(udf->bufSize), .bufLen = udf->bufSize, .numOfResult = 0};
693
      code = udf->scriptPlugin->udfAggStartFunc(&outBuf, udf->scriptUdfCtx);
S
shenglian zhou 已提交
694 695 696 697 698 699 700
      subRsp->resultBuf = outBuf;
      break;
    }
    case TSDB_UDF_CALL_AGG_PROC: {
      SUdfDataBlock input = {0};
      convertDataBlockToUdfDataBlock(&call->block, &input);
      SUdfInterBuf outBuf = {.buf = taosMemoryMalloc(udf->bufSize), .bufLen = udf->bufSize, .numOfResult = 0};
701
      code = udf->scriptPlugin->udfAggProcFunc(&input, &call->interBuf, &outBuf, udf->scriptUdfCtx);
S
shenglian zhou 已提交
702 703 704 705 706 707
      freeUdfInterBuf(&call->interBuf);
      freeUdfDataDataBlock(&input);
      subRsp->resultBuf = outBuf;

      break;
    }
S
slzhou 已提交
708 709
    case TSDB_UDF_CALL_AGG_MERGE: {
      SUdfInterBuf outBuf = {.buf = taosMemoryMalloc(udf->bufSize), .bufLen = udf->bufSize, .numOfResult = 0};
710
      code = udf->scriptPlugin->udfAggMergeFunc(&call->interBuf, &call->interBuf2, &outBuf, udf->scriptUdfCtx);
S
slzhou 已提交
711 712 713 714 715 716
      freeUdfInterBuf(&call->interBuf);
      freeUdfInterBuf(&call->interBuf2);
      subRsp->resultBuf = outBuf;

      break;
    }
S
shenglian zhou 已提交
717 718
    case TSDB_UDF_CALL_AGG_FIN: {
      SUdfInterBuf outBuf = {.buf = taosMemoryMalloc(udf->bufSize), .bufLen = udf->bufSize, .numOfResult = 0};
719
      code = udf->scriptPlugin->udfAggFinishFunc(&call->interBuf, &outBuf, udf->scriptUdfCtx);
S
shenglian zhou 已提交
720 721 722 723 724 725 726 727 728 729
      freeUdfInterBuf(&call->interBuf);
      subRsp->resultBuf = outBuf;
      break;
    }
    default:
      break;
  }

  rsp->seqNum = request->seqNum;
  rsp->type = request->type;
S
slzhou 已提交
730
  rsp->code = (code != 0) ? TSDB_CODE_UDF_FUNC_EXEC_FAILURE : 0;
S
shenglian zhou 已提交
731 732 733 734 735 736 737 738 739 740 741
  subRsp->callType = call->callType;

  int32_t len = encodeUdfResponse(NULL, rsp);
  rsp->msgLen = len;
  void *bufBegin = taosMemoryMalloc(len);
  void *buf = bufBegin;
  encodeUdfResponse(&buf, rsp);
  uvUdf->output = uv_buf_init(bufBegin, len);

  switch (call->callType) {
    case TSDB_UDF_CALL_SCALA_PROC: {
742 743
      blockDataFreeRes(&call->block);
      blockDataFreeRes(&subRsp->resultData);
S
shenglian zhou 已提交
744 745 746 747 748 749 750
      break;
    }
    case TSDB_UDF_CALL_AGG_INIT: {
      freeUdfInterBuf(&subRsp->resultBuf);
      break;
    }
    case TSDB_UDF_CALL_AGG_PROC: {
751
      blockDataFreeRes(&call->block);
S
shenglian zhou 已提交
752 753 754
      freeUdfInterBuf(&subRsp->resultBuf);
      break;
    }
S
slzhou 已提交
755 756 757 758
    case TSDB_UDF_CALL_AGG_MERGE: {
      freeUdfInterBuf(&subRsp->resultBuf);
      break;
    }
S
shenglian zhou 已提交
759 760 761 762 763 764 765 766 767 768 769 770 771 772 773 774
    case TSDB_UDF_CALL_AGG_FIN: {
      freeUdfInterBuf(&subRsp->resultBuf);
      break;
    }
    default:
      break;
  }

  taosMemoryFree(uvUdf->input.base);
  return;
}

void udfdProcessTeardownRequest(SUvUdfWork *uvUdf, SUdfRequest *request) {
  SUdfTeardownRequest *teardown = &request->teardown;
  fnInfo("teardown. seq number: %" PRId64 ", handle:%" PRIx64, request->seqNum, teardown->udfHandle);
  SUdfcFuncHandle *handle = (SUdfcFuncHandle *)(teardown->udfHandle);
775
  SUdf            *udf = handle->udf;
S
shenglian zhou 已提交
776 777 778 779 780
  bool             unloadUdf = false;
  int32_t          code = TSDB_CODE_SUCCESS;

  uv_mutex_lock(&global.udfsMutex);
  udf->refCount--;
S
slzhou 已提交
781
  if (udf->refCount == 0 && (!udf->resident || udf->expired)) {
S
shenglian zhou 已提交
782 783 784 785 786
    unloadUdf = true;
    taosHashRemove(global.udfsHash, udf->name, strlen(udf->name));
  }
  uv_mutex_unlock(&global.udfsMutex);
  if (unloadUdf) {
787
    fnInfo("udf teardown. udf name: %s type %d: context %p", udf->name, udf->scriptType, (void *)(udf->scriptUdfCtx));
S
shenglian zhou 已提交
788 789
    uv_cond_destroy(&udf->condReady);
    uv_mutex_destroy(&udf->lock);
S
slzhou 已提交
790 791
    code = udf->scriptPlugin->udfDestroyFunc(udf->scriptUdfCtx);
    fnDebug("udfd destroy function returns %d", code);
792
    taosRemoveFile(udf->path);
S
slzhou 已提交
793
    taosMemoryFree(udf);
S
shenglian zhou 已提交
794 795 796
  }
  taosMemoryFree(handle);

S
slzhou 已提交
797
  SUdfResponse  response = {0};
S
shenglian zhou 已提交
798 799 800 801 802 803 804 805 806 807 808 809 810 811 812
  SUdfResponse *rsp = &response;
  rsp->seqNum = request->seqNum;
  rsp->type = request->type;
  rsp->code = code;
  int32_t len = encodeUdfResponse(NULL, rsp);
  rsp->msgLen = len;
  void *bufBegin = taosMemoryMalloc(len);
  void *buf = bufBegin;
  encodeUdfResponse(&buf, rsp);
  uvUdf->output = uv_buf_init(bufBegin, len);

  taosMemoryFree(uvUdf->input.base);
  return;
}

S
slzhou 已提交
813
int32_t udfdSaveFuncBodyToFile(SFuncInfo *pFuncInfo, SUdf *udf) {
814 815 816 817 818 819 820 821 822 823 824 825 826 827 828 829 830 831 832 833 834 835 836 837 838 839 840 841
  if (!osTempSpaceAvailable()) {
    terrno = TSDB_CODE_NO_AVAIL_DISK;
    fnError("udfd create shared library failed since %s", terrstr(terrno));
    return terrno;
  }

  char path[PATH_MAX] = {0};
#ifdef WINDOWS
  snprintf(path, sizeof(path), "%s%s_%d_%" PRId64, tsTempDir, pFuncInfo->name, udf->version, udf->lastFetchTime);
#else
  snprintf(path, sizeof(path), "%s/%s_%d_%" PRId64, tsTempDir, pFuncInfo->name, udf->version, udf->lastFetchTime);
#endif
  TdFilePtr file = taosOpenFile(path, TD_FILE_CREATE | TD_FILE_WRITE | TD_FILE_READ | TD_FILE_TRUNC);
  if (file == NULL) {
    fnError("udfd write udf shared library: %s failed, error: %d %s", path, errno, strerror(errno));
    return TSDB_CODE_FILE_CORRUPTED;
  }

  int64_t count = taosWriteFile(file, pFuncInfo->pCode, pFuncInfo->codeSize);
  if (count != pFuncInfo->codeSize) {
    fnError("udfd write udf shared library failed");
    return TSDB_CODE_FILE_CORRUPTED;
  }
  taosCloseFile(&file);
  strncpy(udf->path, path, PATH_MAX);
  return TSDB_CODE_SUCCESS;
}

842
void udfdProcessRpcRsp(void *parent, SRpcMsg *pMsg, SEpSet *pEpSet) {
843
  SUdfdRpcSendRecvInfo *msgInfo = (SUdfdRpcSendRecvInfo *)pMsg->info.ahandle;
844 845 846 847 848 849 850 851 852 853 854 855 856 857 858 859

  if (pEpSet) {
    if (!isEpsetEqual(&global.mgmtEp.epSet, pEpSet)) {
      updateEpSet_s(&global.mgmtEp, pEpSet);
    }
  }

  if (pMsg->code != TSDB_CODE_SUCCESS) {
    fnError("udfd rpc error. code: %s", tstrerror(pMsg->code));
    msgInfo->code = pMsg->code;
    goto _return;
  }

  if (msgInfo->rpcType == UDFD_RPC_MNODE_CONNECT) {
    SConnectRsp connectRsp = {0};
    tDeserializeSConnectRsp(pMsg->pCont, pMsg->contLen, &connectRsp);
860

dengyihao's avatar
dengyihao 已提交
861 862 863 864 865 866
    int32_t now = taosGetTimestampSec();
    int32_t delta = abs(now - connectRsp.svrTimestamp);
    if (delta > 900) {
      msgInfo->code = TSDB_CODE_TIME_UNSYNCED;
      goto _return;
    }
867

868
    if (connectRsp.epSet.numOfEps == 0) {
S
Shengliang Guan 已提交
869
      msgInfo->code = TSDB_CODE_APP_ERROR;
870 871 872 873 874 875 876 877 878 879
      goto _return;
    }

    if (connectRsp.dnodeNum > 1 && !isEpsetEqual(&global.mgmtEp.epSet, &connectRsp.epSet)) {
      updateEpSet_s(&global.mgmtEp, &connectRsp.epSet);
    }
    msgInfo->code = 0;
  } else if (msgInfo->rpcType == UDFD_RPC_RETRIVE_FUNC) {
    SRetrieveFuncRsp retrieveRsp = {0};
    tDeserializeSRetrieveFuncRsp(pMsg->pCont, pMsg->contLen, &retrieveRsp);
880

881
    SFuncInfo *pFuncInfo = (SFuncInfo *)taosArrayGet(retrieveRsp.pFuncInfos, 0);
882
    SUdf      *udf = msgInfo->param;
883 884
    udf->funcType = pFuncInfo->funcType;
    udf->scriptType = pFuncInfo->scriptType;
885
    udf->outputType = pFuncInfo->outputType;
886 887
    udf->outputLen = pFuncInfo->outputLen;
    udf->bufSize = pFuncInfo->bufSize;
888
    udf->version = *(int32_t *)taosArrayGet(retrieveRsp.pFuncVersions, 0);
889

890 891 892
    msgInfo->code = udfdSaveFuncBodyToFile(pFuncInfo, udf);
    if (msgInfo->code == 0) {
      udf->lastFetchTime = taosGetTimestampMs();
S
shenglian zhou 已提交
893
    }
894 895
    tFreeSFuncInfo(pFuncInfo);
    taosArrayDestroy(retrieveRsp.pFuncInfos);
896
    taosArrayDestroy(retrieveRsp.pFuncVersions);
897 898 899 900 901 902 903 904 905 906 907 908 909 910 911
  }

_return:
  rpcFreeCont(pMsg->pCont);
  uv_sem_post(&msgInfo->resultSem);
  return;
}

int32_t udfdFillUdfInfoFromMNode(void *clientRpc, char *udfName, SUdf *udf) {
  SRetrieveFuncReq retrieveReq = {0};
  retrieveReq.numOfFuncs = 1;
  retrieveReq.pFuncNames = taosArrayInit(1, TSDB_FUNC_NAME_LEN);
  taosArrayPush(retrieveReq.pFuncNames, udfName);

  int32_t contLen = tSerializeSRetrieveFuncReq(NULL, 0, &retrieveReq);
912
  void   *pReq = rpcMallocCont(contLen);
913 914 915
  tSerializeSRetrieveFuncReq(pReq, contLen, &retrieveReq);
  taosArrayDestroy(retrieveReq.pFuncNames);

dengyihao's avatar
dengyihao 已提交
916
  SUdfdRpcSendRecvInfo *msgInfo = taosMemoryCalloc(1, sizeof(SUdfdRpcSendRecvInfo));
917 918 919 920 921 922 923 924
  msgInfo->rpcType = UDFD_RPC_RETRIVE_FUNC;
  msgInfo->param = udf;
  uv_sem_init(&msgInfo->resultSem, 0);

  SRpcMsg rpcMsg = {0};
  rpcMsg.pCont = pReq;
  rpcMsg.contLen = contLen;
  rpcMsg.msgType = TDMT_MND_RETRIEVE_FUNC;
925
  rpcMsg.info.ahandle = msgInfo;
926 927 928 929 930 931 932 933 934 935 936 937
  rpcSendRequest(clientRpc, &global.mgmtEp.epSet, &rpcMsg, NULL);

  uv_sem_wait(&msgInfo->resultSem);
  uv_sem_destroy(&msgInfo->resultSem);
  int32_t code = msgInfo->code;
  taosMemoryFree(msgInfo);
  return code;
}

int32_t udfdConnectToMnode() {
  SConnectReq connReq = {0};
  connReq.connType = CONN_TYPE__UDFD;
dengyihao's avatar
dengyihao 已提交
938
  tstrncpy(connReq.app, "udfd", sizeof(connReq.app));
939 940 941 942
  tstrncpy(connReq.user, TSDB_DEFAULT_USER, sizeof(connReq.user));
  char pass[TSDB_PASSWORD_LEN + 1] = {0};
  taosEncryptPass_c((uint8_t *)(TSDB_DEFAULT_PASS), strlen(TSDB_DEFAULT_PASS), pass);
  tstrncpy(connReq.passwd, pass, sizeof(connReq.passwd));
943 944
  connReq.pid = taosGetPId();
  connReq.startTime = taosGetTimestampMs();
D
dapan1121 已提交
945
  strcpy(connReq.sVer, version);
946 947

  int32_t contLen = tSerializeSConnectReq(NULL, 0, &connReq);
948
  void   *pReq = rpcMallocCont(contLen);
949 950 951 952 953 954 955 956 957 958
  tSerializeSConnectReq(pReq, contLen, &connReq);

  SUdfdRpcSendRecvInfo *msgInfo = taosMemoryCalloc(1, sizeof(SUdfdRpcSendRecvInfo));
  msgInfo->rpcType = UDFD_RPC_MNODE_CONNECT;
  uv_sem_init(&msgInfo->resultSem, 0);

  SRpcMsg rpcMsg = {0};
  rpcMsg.msgType = TDMT_MND_CONNECT;
  rpcMsg.pCont = pReq;
  rpcMsg.contLen = contLen;
959
  rpcMsg.info.ahandle = msgInfo;
960 961 962 963 964 965 966 967
  rpcSendRequest(global.clientRpc, &global.mgmtEp.epSet, &rpcMsg, NULL);

  uv_sem_wait(&msgInfo->resultSem);
  int32_t code = msgInfo->code;
  uv_sem_destroy(&msgInfo->resultSem);
  taosMemoryFree(msgInfo);
  return code;
}
968

dengyihao's avatar
dengyihao 已提交
969
static bool udfdRpcRfp(int32_t code, tmsg_t msgType) {
970
  if (code == TSDB_CODE_RPC_NETWORK_UNAVAIL || code == TSDB_CODE_RPC_BROKEN_LINK || code == TSDB_CODE_SYN_NOT_LEADER ||
971 972
      code == TSDB_CODE_RPC_SOMENODE_NOT_CONNECTED || code == TSDB_CODE_SYN_RESTORING ||
      code == TSDB_CODE_MNODE_NOT_FOUND || code == TSDB_CODE_APP_IS_STARTING || code == TSDB_CODE_APP_IS_STOPPING) {
973 974
    if (msgType == TDMT_SCH_QUERY || msgType == TDMT_SCH_MERGE_QUERY || msgType == TDMT_SCH_FETCH ||
        msgType == TDMT_SCH_MERGE_FETCH) {
dengyihao's avatar
dengyihao 已提交
975
      return false;
976
    }
S
shenglian zhou 已提交
977 978 979 980 981 982 983 984 985 986 987 988 989 990 991 992 993 994 995 996 997 998 999 1000 1001 1002 1003 1004 1005 1006 1007 1008 1009 1010 1011 1012 1013 1014 1015 1016 1017 1018 1019 1020 1021 1022 1023 1024 1025 1026 1027 1028 1029 1030 1031 1032 1033 1034
    return true;
  } else {
    return false;
  }
}

int initEpSetFromCfg(const char *firstEp, const char *secondEp, SCorEpSet *pEpSet) {
  pEpSet->version = 0;

  // init mnode ip set
  SEpSet *mgmtEpSet = &(pEpSet->epSet);
  mgmtEpSet->numOfEps = 0;
  mgmtEpSet->inUse = 0;

  if (firstEp && firstEp[0] != 0) {
    if (strlen(firstEp) >= TSDB_EP_LEN) {
      terrno = TSDB_CODE_TSC_INVALID_FQDN;
      return -1;
    }

    int32_t code = taosGetFqdnPortFromEp(firstEp, &mgmtEpSet->eps[0]);
    if (code != TSDB_CODE_SUCCESS) {
      terrno = TSDB_CODE_TSC_INVALID_FQDN;
      return terrno;
    }

    mgmtEpSet->numOfEps++;
  }

  if (secondEp && secondEp[0] != 0) {
    if (strlen(secondEp) >= TSDB_EP_LEN) {
      terrno = TSDB_CODE_TSC_INVALID_FQDN;
      return -1;
    }

    taosGetFqdnPortFromEp(secondEp, &mgmtEpSet->eps[mgmtEpSet->numOfEps]);
    mgmtEpSet->numOfEps++;
  }

  if (mgmtEpSet->numOfEps == 0) {
    terrno = TSDB_CODE_TSC_INVALID_FQDN;
    return -1;
  }

  return 0;
}

int32_t udfdOpenClientRpc() {
  SRpcInit rpcInit = {0};
  rpcInit.label = "UDFD";
  rpcInit.numOfThreads = 1;
  rpcInit.cfp = (RpcCfp)udfdProcessRpcRsp;
  rpcInit.sessions = 1024;
  rpcInit.connType = TAOS_CONN_CLIENT;
  rpcInit.idleTime = tsShellActivityTimer * 1000;
  rpcInit.user = TSDB_DEFAULT_USER;
  rpcInit.parent = &global;
  rpcInit.rfp = udfdRpcRfp;
dengyihao's avatar
dengyihao 已提交
1035
  rpcInit.compressSize = tsCompressMsgSize;
1036

dengyihao's avatar
dengyihao 已提交
1037 1038 1039 1040 1041 1042
  int32_t connLimitNum = tsNumOfRpcSessions / (tsNumOfRpcThreads * 3);
  connLimitNum = TMAX(connLimitNum, 10);
  connLimitNum = TMIN(connLimitNum, 500);
  rpcInit.connLimitNum = connLimitNum;
  rpcInit.timeToGetConn = tsTimeToGetAvailableConn;

S
shenglian zhou 已提交
1043 1044 1045 1046 1047 1048 1049 1050 1051
  global.clientRpc = rpcOpen(&rpcInit);
  if (global.clientRpc == NULL) {
    fnError("failed to init dnode rpc client");
    return -1;
  }
  return 0;
}

int32_t udfdCloseClientRpc() {
1052
  fnInfo("udfd begin closing rpc");
S
shenglian zhou 已提交
1053
  rpcClose(global.clientRpc);
1054
  fnInfo("udfd finish closing rpc");
S
shenglian zhou 已提交
1055 1056 1057
  return 0;
}

1058
void udfdOnWrite(uv_write_t *req, int status) {
S
slzhou 已提交
1059 1060
  SUvUdfWork *work = (SUvUdfWork *)req->data;
  if (status < 0) {
S
slzhou 已提交
1061
    fnError("udfd send response error, length: %zu code: %s", work->output.len, uv_err_name(status));
S
slzhou 已提交
1062
  }
1063 1064 1065 1066 1067 1068 1069 1070 1071 1072 1073
  // remove work from the connection work list
  if (work->conn != NULL) {
    SUvUdfWork **ppWork;
    for (ppWork = &work->conn->pWorkList; *ppWork && (*ppWork != work); ppWork = &((*ppWork)->pWorkNext)) {
    }
    if (*ppWork == work) {
      *ppWork = work->pWorkNext;
    } else {
      fnError("work not in conn any more");
    }
  }
S
slzhou 已提交
1074 1075 1076
  taosMemoryFree(work->output.base);
  taosMemoryFree(work);
  taosMemoryFree(req);
1077 1078 1079
}

void udfdSendResponse(uv_work_t *work, int status) {
S
slzhou 已提交
1080
  SUvUdfWork *udfWork = (SUvUdfWork *)(work->data);
1081

1082 1083 1084 1085 1086
  if (udfWork->conn != NULL) {
    uv_write_t *write_req = taosMemoryMalloc(sizeof(uv_write_t));
    write_req->data = udfWork;
    uv_write(write_req, udfWork->conn->client, &udfWork->output, 1, udfdOnWrite);
  }
S
slzhou 已提交
1087
  taosMemoryFree(work);
1088 1089 1090
}

void udfdAllocBuffer(uv_handle_t *handle, size_t suggestedSize, uv_buf_t *buf) {
S
slzhou 已提交
1091 1092 1093 1094 1095 1096 1097 1098 1099 1100 1101
  SUdfdUvConn *ctx = handle->data;
  int32_t      msgHeadSize = sizeof(int32_t) + sizeof(int64_t);
  if (ctx->inputCap == 0) {
    ctx->inputBuf = taosMemoryMalloc(msgHeadSize);
    if (ctx->inputBuf) {
      ctx->inputLen = 0;
      ctx->inputCap = msgHeadSize;
      ctx->inputTotal = -1;

      buf->base = ctx->inputBuf;
      buf->len = ctx->inputCap;
1102
    } else {
dengyihao's avatar
dengyihao 已提交
1103
      fnError("udfd can not allocate enough memory") buf->base = NULL;
S
slzhou 已提交
1104
      buf->len = 0;
1105
    }
1106
  } else if (ctx->inputTotal == -1 && ctx->inputLen < msgHeadSize) {
1107 1108
    buf->base = ctx->inputBuf + ctx->inputLen;
    buf->len = msgHeadSize - ctx->inputLen;
S
slzhou 已提交
1109 1110 1111 1112 1113 1114 1115 1116
  } else {
    ctx->inputCap = ctx->inputTotal > ctx->inputCap ? ctx->inputTotal : ctx->inputCap;
    void *inputBuf = taosMemoryRealloc(ctx->inputBuf, ctx->inputCap);
    if (inputBuf) {
      ctx->inputBuf = inputBuf;
      buf->base = ctx->inputBuf + ctx->inputLen;
      buf->len = ctx->inputCap - ctx->inputLen;
    } else {
dengyihao's avatar
dengyihao 已提交
1117
      fnError("udfd can not allocate enough memory") buf->base = NULL;
S
slzhou 已提交
1118 1119 1120
      buf->len = 0;
    }
  }
1121 1122 1123
}

bool isUdfdUvMsgComplete(SUdfdUvConn *pipe) {
S
slzhou 已提交
1124 1125 1126 1127 1128 1129 1130 1131
  if (pipe->inputTotal == -1 && pipe->inputLen >= sizeof(int32_t)) {
    pipe->inputTotal = *(int32_t *)(pipe->inputBuf);
  }
  if (pipe->inputLen == pipe->inputCap && pipe->inputTotal == pipe->inputCap) {
    fnDebug("receive request complete. length %d", pipe->inputLen);
    return true;
  }
  return false;
1132 1133 1134
}

void udfdHandleRequest(SUdfdUvConn *conn) {
1135
  char   *inputBuf = conn->inputBuf;
S
shenglian zhou 已提交
1136 1137
  int32_t inputLen = conn->inputLen;

1138
  uv_work_t  *work = taosMemoryMalloc(sizeof(uv_work_t));
S
slzhou 已提交
1139
  SUvUdfWork *udfWork = taosMemoryMalloc(sizeof(SUvUdfWork));
1140 1141 1142
  udfWork->conn = conn;
  udfWork->pWorkNext = conn->pWorkList;
  conn->pWorkList = udfWork;
S
shenglian zhou 已提交
1143
  udfWork->input = uv_buf_init(inputBuf, inputLen);
S
slzhou 已提交
1144 1145 1146 1147 1148 1149
  conn->inputBuf = NULL;
  conn->inputLen = 0;
  conn->inputCap = 0;
  conn->inputTotal = -1;
  work->data = udfWork;
  uv_queue_work(global.loop, work, udfdProcessRequest, udfdSendResponse);
1150 1151 1152
}

void udfdPipeCloseCb(uv_handle_t *pipe) {
S
slzhou 已提交
1153
  SUdfdUvConn *conn = pipe->data;
1154
  SUvUdfWork  *pWork = conn->pWorkList;
1155 1156 1157 1158 1159
  while (pWork != NULL) {
    pWork->conn = NULL;
    pWork = pWork->pWorkNext;
  }

S
slzhou 已提交
1160 1161 1162
  taosMemoryFree(conn->client);
  taosMemoryFree(conn->inputBuf);
  taosMemoryFree(conn);
1163 1164 1165
}

void udfdPipeRead(uv_stream_t *client, ssize_t nread, const uv_buf_t *buf) {
1166
  fnDebug("udfd read %zd bytes from client", nread);
S
slzhou 已提交
1167
  if (nread == 0) return;
1168

S
slzhou 已提交
1169
  SUdfdUvConn *conn = client->data;
1170

S
slzhou 已提交
1171 1172 1173 1174 1175 1176
  if (nread > 0) {
    conn->inputLen += nread;
    if (isUdfdUvMsgComplete(conn)) {
      udfdHandleRequest(conn);
    } else {
      // log error or continue;
1177
    }
S
slzhou 已提交
1178 1179
    return;
  }
1180

S
slzhou 已提交
1181 1182
  if (nread < 0) {
    if (nread == UV_EOF) {
S
slzhou 已提交
1183
      fnInfo("udfd pipe read EOF");
S
slzhou 已提交
1184
    } else {
1185
      fnError("Receive error %s", uv_err_name(nread));
1186
    }
S
slzhou 已提交
1187 1188
    udfdUvHandleError(conn);
  }
1189 1190 1191
}

void udfdOnNewConnection(uv_stream_t *server, int status) {
S
slzhou 已提交
1192
  if (status < 0) {
S
slzhou 已提交
1193
    fnError("udfd new connection error. code: %s", uv_strerror(status));
S
slzhou 已提交
1194 1195
    return;
  }
1196

S
slzhou 已提交
1197 1198 1199 1200
  uv_pipe_t *client = (uv_pipe_t *)taosMemoryMalloc(sizeof(uv_pipe_t));
  uv_pipe_init(global.loop, client, 0);
  if (uv_accept(server, (uv_stream_t *)client) == 0) {
    SUdfdUvConn *ctx = taosMemoryMalloc(sizeof(SUdfdUvConn));
1201
    ctx->pWorkList = NULL;
S
slzhou 已提交
1202 1203 1204 1205 1206 1207 1208 1209 1210 1211
    ctx->client = (uv_stream_t *)client;
    ctx->inputBuf = 0;
    ctx->inputLen = 0;
    ctx->inputCap = 0;
    client->data = ctx;
    ctx->client = (uv_stream_t *)client;
    uv_read_start((uv_stream_t *)client, udfdAllocBuffer, udfdPipeRead);
  } else {
    uv_close((uv_handle_t *)client, NULL);
  }
1212 1213
}

1214 1215
void udfdIntrSignalHandler(uv_signal_t *handle, int signum) {
  fnInfo("udfd signal received: %d\n", signum);
S
slzhou 已提交
1216
  uv_fs_t req;
1217 1218 1219
  uv_fs_unlink(global.loop, &req, global.listenPipeName, NULL);
  uv_signal_stop(handle);
  uv_stop(global.loop);
1220 1221
}

S
slzhou 已提交
1222 1223 1224 1225 1226 1227 1228 1229 1230 1231 1232 1233 1234 1235 1236 1237 1238 1239
static int32_t udfdParseArgs(int32_t argc, char *argv[]) {
  for (int32_t i = 1; i < argc; ++i) {
    if (strcmp(argv[i], "-c") == 0) {
      if (i < argc - 1) {
        if (strlen(argv[++i]) >= PATH_MAX) {
          printf("config file path overflow");
          return -1;
        }
        tstrncpy(configDir, argv[i], PATH_MAX);
      } else {
        printf("'-c' requires a parameter, default is %s\n", configDir);
        return -1;
      }
    } else if (strcmp(argv[i], "-V") == 0) {
      global.printVersion = true;
    } else {
    }
  }
1240 1241 1242 1243

  return 0;
}

S
shenglian zhou 已提交
1244 1245 1246 1247 1248 1249 1250 1251 1252 1253 1254
static void udfdPrintVersion() {
#ifdef TD_ENTERPRISE
  char *releaseName = "enterprise";
#else
  char *releaseName = "community";
#endif
  printf("%s version: %s compatible_version: %s\n", releaseName, version, compatible_version);
  printf("gitinfo: %s\n", gitinfo);
  printf("buildInfo: %s\n", buildinfo);
}

S
slzhou 已提交
1255 1256 1257
static int32_t udfdInitLog() {
  char logName[12] = {0};
  snprintf(logName, sizeof(logName), "%slog", "udfd");
wafwerar's avatar
wafwerar 已提交
1258
  return taosCreateLog(logName, 1, configDir, NULL, NULL, NULL, NULL, 0);
S
slzhou 已提交
1259
}
1260

1261 1262 1263 1264 1265 1266 1267 1268
void udfdCtrlAllocBufCb(uv_handle_t *handle, size_t suggested_size, uv_buf_t *buf) {
  buf->base = taosMemoryMalloc(suggested_size);
  buf->len = suggested_size;
}

void udfdCtrlReadCb(uv_stream_t *q, ssize_t nread, const uv_buf_t *buf) {
  if (nread < 0) {
    fnError("udfd ctrl pipe read error. %s", uv_err_name(nread));
1269
    taosMemoryFree(buf->base);
S
slzhou 已提交
1270
    uv_close((uv_handle_t *)q, NULL);
1271 1272 1273 1274 1275 1276 1277 1278 1279
    uv_stop(global.loop);
    return;
  }
  fnError("udfd ctrl pipe read %zu bytes", nread);
  taosMemoryFree(buf->base);
}

static int32_t removeListeningPipe() {
  uv_fs_t req;
S
slzhou 已提交
1280
  int     err = uv_fs_unlink(global.loop, &req, global.listenPipeName, NULL);
1281 1282 1283 1284
  uv_fs_req_cleanup(&req);
  return err;
}

S
slzhou 已提交
1285
static int32_t udfdUvInit() {
S
slzhou 已提交
1286
  uv_loop_t *loop = taosMemoryMalloc(sizeof(uv_loop_t));
S
slzhou 已提交
1287 1288
  if (loop) {
    uv_loop_init(loop);
S
slzhou 已提交
1289 1290
  } else {
    return -1;
S
slzhou 已提交
1291 1292
  }
  global.loop = loop;
1293

1294
  if (tsStartUdfd) {  // udfd is started by taosd, which shall exit when taosd exit
1295 1296 1297 1298
    uv_pipe_init(global.loop, &global.ctrlPipe, 1);
    uv_pipe_open(&global.ctrlPipe, 0);
    uv_read_start((uv_stream_t *)&global.ctrlPipe, udfdCtrlAllocBufCb, udfdCtrlReadCb);
  }
1299
  getUdfdPipeName(global.listenPipeName, sizeof(global.listenPipeName));
S
slzhou 已提交
1300

1301
  removeListeningPipe();
S
slzhou 已提交
1302

1303
  uv_pipe_init(global.loop, &global.listeningPipe, 0);
S
slzhou 已提交
1304

1305 1306
  uv_signal_init(global.loop, &global.intrSignal);
  uv_signal_start(&global.intrSignal, udfdIntrSignalHandler, SIGINT);
S
slzhou 已提交
1307 1308 1309

  int r;
  fnInfo("bind to pipe %s", global.listenPipeName);
1310
  if ((r = uv_pipe_bind(&global.listeningPipe, global.listenPipeName))) {
S
slzhou 已提交
1311
    fnError("Bind error %s", uv_err_name(r));
1312
    removeListeningPipe();
S
slzhou 已提交
1313
    return -2;
S
slzhou 已提交
1314
  }
1315
  if ((r = uv_listen((uv_stream_t *)&global.listeningPipe, 128, udfdOnNewConnection))) {
S
slzhou 已提交
1316
    fnError("Listen error %s", uv_err_name(r));
1317
    removeListeningPipe();
S
slzhou 已提交
1318
    return -3;
S
slzhou 已提交
1319
  }
1320 1321
  return 0;
}
1322

dengyihao's avatar
dengyihao 已提交
1323
static void udfdCloseWalkCb(uv_handle_t *handle, void *arg) {
S
slzhou 已提交
1324 1325 1326 1327 1328
  if (!uv_is_closing(handle)) {
    uv_close(handle, NULL);
  }
}

S
slzhou 已提交
1329
static int32_t udfdRun() {
1330 1331
  uv_mutex_init(&global.scriptPluginsMutex);

S
slzhou 已提交
1332 1333
  global.udfsHash = taosHashInit(64, taosGetDefaultHashFunction(TSDB_DATA_TYPE_BINARY), true, HASH_NO_LOCK);
  uv_mutex_init(&global.udfsMutex);
1334

S
slzhou 已提交
1335 1336 1337 1338 1339 1340 1341 1342 1343 1344
  fnInfo("start udfd event loop");
  uv_run(global.loop, UV_RUN_DEFAULT);
  fnInfo("udfd event loop stopped.");

  uv_loop_close(global.loop);

  uv_walk(global.loop, udfdCloseWalkCb, NULL);
  uv_run(global.loop, UV_RUN_DEFAULT);
  uv_loop_close(global.loop);

S
slzhou 已提交
1345
  return 0;
S
slzhou 已提交
1346
}
1347

dengyihao's avatar
dengyihao 已提交
1348
void udfdConnectMnodeThreadFunc(void *args) {
1349 1350 1351 1352 1353 1354 1355 1356 1357 1358 1359 1360 1361 1362 1363 1364
  int32_t retryMnodeTimes = 0;
  int32_t code = 0;
  while (retryMnodeTimes++ <= TSDB_MAX_REPLICA) {
    uv_sleep(100 * (1 << retryMnodeTimes));
    code = udfdConnectToMnode();
    if (code == 0) {
      break;
    }
    fnError("udfd can not connect to mnode, code: %s. retry", tstrerror(code));
  }

  if (code != 0) {
    fnError("udfd can not connect to mnode");
  }
}

1365
int32_t udfdInitResidentFuncs() {
S
slzhou 已提交
1366 1367 1368 1369
  if (strlen(tsUdfdResFuncs) == 0) {
    return TSDB_CODE_SUCCESS;
  }

1370
  global.residentFuncs = taosArrayInit(2, TSDB_FUNC_NAME_LEN);
1371 1372
  char *pSave = tsUdfdResFuncs;
  char *token;
S
slzhou 已提交
1373
  while ((token = strtok_r(pSave, ",", &pSave)) != NULL) {
1374
    char func[TSDB_FUNC_NAME_LEN + 1] = {0};
S
shenglian zhou 已提交
1375
    strncpy(func, token, TSDB_FUNC_NAME_LEN);
1376
    fnInfo("udfd add resident function %s", func);
S
slzhou 已提交
1377 1378 1379
    taosArrayPush(global.residentFuncs, func);
  }

1380 1381 1382 1383
  return TSDB_CODE_SUCCESS;
}

int32_t udfdDeinitResidentFuncs() {
1384
  for (int32_t i = 0; i < taosArrayGetSize(global.residentFuncs); ++i) {
1385 1386
    char  *funcName = taosArrayGet(global.residentFuncs, i);
    SUdf **udfInHash = taosHashGet(global.udfsHash, funcName, strlen(funcName));
S
slzhou 已提交
1387
    if (udfInHash) {
1388
      SUdf   *udf = *udfInHash;
S
slzhou 已提交
1389 1390
      int32_t code = udf->scriptPlugin->udfDestroyFunc(udf->scriptUdfCtx);
      fnDebug("udfd destroy function returns %d", code);
S
slzhou 已提交
1391
      taosHashRemove(global.udfsHash, funcName, strlen(funcName));
S
slzhou 已提交
1392
      taosMemoryFree(udf);
S
slzhou 已提交
1393 1394
    }
  }
1395
  taosArrayDestroy(global.residentFuncs);
1396 1397 1398
  return TSDB_CODE_SUCCESS;
}

1399 1400 1401 1402 1403 1404
int32_t udfdCleanup() {
  uv_mutex_destroy(&global.udfsMutex);
  taosHashCleanup(global.udfsHash);
  return 0;
}

S
slzhou 已提交
1405
int main(int argc, char *argv[]) {
D
dapan1121 已提交
1406 1407
  if (!taosCheckSystemIsLittleEnd()) {
    printf("failed to start since on non-little-end machines\n");
S
slzhou 已提交
1408 1409 1410 1411 1412 1413 1414 1415 1416 1417 1418 1419 1420 1421
    return -1;
  }

  if (udfdParseArgs(argc, argv) != 0) {
    printf("failed to start since parse args error\n");
    return -1;
  }

  if (global.printVersion) {
    udfdPrintVersion();
    return 0;
  }

  if (udfdInitLog() != 0) {
A
Alex Duan 已提交
1422
    // ignore create log failed, because this error no matter
S
slzhou 已提交
1423 1424 1425
    printf("failed to start since init log error\n");
  }

wafwerar's avatar
wafwerar 已提交
1426
  if (taosInitCfg(configDir, NULL, NULL, NULL, NULL, 0) != 0) {
S
slzhou 已提交
1427
    fnError("failed to start since read config error");
1428
    return -2;
S
slzhou 已提交
1429 1430
  }

1431
  initEpSetFromCfg(tsFirst, tsSecond, &global.mgmtEp);
1432 1433 1434 1435 1436
  if (udfdOpenClientRpc() != 0) {
    fnError("open rpc connection to mnode failure");
    return -3;
  }

S
slzhou 已提交
1437 1438 1439 1440 1441
  if (udfdUvInit() != 0) {
    fnError("uv init failure");
    return -5;
  }

1442 1443
  udfdInitResidentFuncs();

1444 1445 1446
  uv_thread_t mnodeConnectThread;
  uv_thread_create(&mnodeConnectThread, udfdConnectMnodeThreadFunc, NULL);

1447
  udfdRun();
dengyihao's avatar
dengyihao 已提交
1448

S
slzhou 已提交
1449 1450
  removeListeningPipe();
  udfdCloseClientRpc();
S
shenglian zhou 已提交
1451

1452
  udfdDeinitResidentFuncs();
1453 1454 1455

  udfdDeinitScriptPlugins();

1456
  udfdCleanup();
S
shenglian zhou 已提交
1457
  return 0;
1458
}