tudf.c 48.9 KB
Newer Older
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16
/*
 * 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/>.
 */
#include "uv.h"
#include "os.h"
S
slzhou 已提交
17
#include "fnLog.h"
18
#include "tudf.h"
19
#include "tudfInt.h"
S
shenglian zhou 已提交
20
#include "tarray.h"
S
slzhou 已提交
21
#include "tglobal.h"
S
slzhou 已提交
22
#include "tdatablock.h"
S
shenglian zhou 已提交
23 24
#include "querynodes.h"
#include "builtinsimpl.h"
S
slzhou 已提交
25
#include "functionMgt.h"
26

27
//TODO: add unit test
S
shenglian zhou 已提交
28
//TODO: include all global variable under context struct
S
slzhou 已提交
29

30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141
typedef struct SUdfdData {
  bool          startCalled;
  bool          needCleanUp;
  uv_loop_t     loop;
  uv_thread_t   thread;
  uv_barrier_t  barrier;
  uv_process_t  process;
  int           spawnErr;
  uv_pipe_t     ctrlPipe;
  uv_async_t    stopAsync;
  int32_t        stopCalled;

  int32_t         dnodeId;
} SUdfdData;

SUdfdData udfdGlobal = {0};

static int32_t udfSpawnUdfd(SUdfdData *pData);

void udfUdfdExit(uv_process_t *process, int64_t exitStatus, int termSignal) {
  fnInfo("udfd process exited with status %" PRId64 ", signal %d", exitStatus, termSignal);
  SUdfdData *pData = process->data;
  if (exitStatus == 0 && termSignal == 0 || atomic_load_32(&pData->stopCalled)) {
    fnInfo("udfd process exit due to SIGINT or dnode-mgmt called stop");
  } else {
    fnInfo("udfd process restart");
    udfSpawnUdfd(pData);
  }
}

static int32_t udfSpawnUdfd(SUdfdData* pData) {
  fnInfo("dnode start spawning udfd");
  uv_process_options_t options = {0};

  char path[PATH_MAX] = {0};
  if (tsProcPath == NULL) {
    path[0] = '.';
  } else {
    strncpy(path, tsProcPath, strlen(tsProcPath));
    taosDirName(path);
  }
#ifdef WINDOWS
  strcat(path, "udfd.exe");
#else
  strcat(path, "/udfd");
#endif
  char* argsUdfd[] = {path, "-c", configDir, NULL};
  options.args = argsUdfd;
  options.file = path;

  options.exit_cb = udfUdfdExit;

  uv_pipe_init(&pData->loop, &pData->ctrlPipe, 1);

  uv_stdio_container_t child_stdio[3];
  child_stdio[0].flags = UV_CREATE_PIPE | UV_READABLE_PIPE;
  child_stdio[0].data.stream = (uv_stream_t*) &pData->ctrlPipe;
  child_stdio[1].flags = UV_IGNORE;
  child_stdio[2].flags = UV_INHERIT_FD;
  child_stdio[2].data.fd = 2;
  options.stdio_count = 3;
  options.stdio = child_stdio;

  options.flags = UV_PROCESS_DETACHED;

  char dnodeIdEnvItem[32] = {0};
  char thrdPoolSizeEnvItem[32] = {0};
  snprintf(dnodeIdEnvItem, 32, "%s=%d", "DNODE_ID", pData->dnodeId);
  float numCpuCores = 4;
  taosGetCpuCores(&numCpuCores);
  snprintf(thrdPoolSizeEnvItem,32,  "%s=%d", "UV_THREADPOOL_SIZE", (int)numCpuCores*2);
  char* envUdfd[] = {dnodeIdEnvItem, thrdPoolSizeEnvItem, NULL};
  options.env = envUdfd;

  int err = uv_spawn(&pData->loop, &pData->process, &options);
  pData->process.data = (void*)pData;

  if (err != 0) {
    fnError("can not spawn udfd. path: %s, error: %s", path, uv_strerror(err));
  }
  return err;
}

static void udfUdfdCloseWalkCb(uv_handle_t* handle, void* arg) {
  if (!uv_is_closing(handle)) {
    uv_close(handle, NULL);
  }
}

static void udfUdfdStopAsyncCb(uv_async_t *async) {
  SUdfdData *pData = async->data;
  uv_stop(&pData->loop);
}

static void udfWatchUdfd(void *args) {
  SUdfdData *pData = args;
  uv_loop_init(&pData->loop);
  uv_async_init(&pData->loop, &pData->stopAsync, udfUdfdStopAsyncCb);
  pData->stopAsync.data = pData;
  int32_t err = udfSpawnUdfd(pData);
  atomic_store_32(&pData->spawnErr, err);
  uv_barrier_wait(&pData->barrier);
  uv_run(&pData->loop, UV_RUN_DEFAULT);
  uv_loop_close(&pData->loop);

  uv_walk(&pData->loop, udfUdfdCloseWalkCb, NULL);
  uv_run(&pData->loop, UV_RUN_DEFAULT);
  uv_loop_close(&pData->loop);
  return;
}

int32_t udfStartUdfd(int32_t startDnodeId) {
S
slzhou 已提交
142 143 144
  if (!tsStartUdfd) {
    fnInfo("start udfd is disabled.")
  }
145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188
  SUdfdData *pData = &udfdGlobal;
  if (pData->startCalled) {
    fnInfo("dnode-mgmt start udfd already called");
    return 0;
  }
  pData->startCalled = true;
  char dnodeId[8] = {0};
  snprintf(dnodeId, sizeof(dnodeId), "%d", startDnodeId);
  uv_os_setenv("DNODE_ID", dnodeId);
  pData->dnodeId = startDnodeId;

  uv_barrier_init(&pData->barrier, 2);
  uv_thread_create(&pData->thread, udfWatchUdfd, pData);
  uv_barrier_wait(&pData->barrier);
  int32_t err = atomic_load_32(&pData->spawnErr);
  if (err != 0) {
    uv_barrier_destroy(&pData->barrier);
    uv_async_send(&pData->stopAsync);
    uv_thread_join(&pData->thread);
    pData->needCleanUp = false;
    fnInfo("dnode-mgmt udfd cleaned up after spawn err");
  } else {
    pData->needCleanUp = true;
  }
  return err;
}

int32_t udfStopUdfd() {
  SUdfdData *pData = &udfdGlobal;
  fnInfo("dnode-mgmt to stop udfd. need cleanup: %d, spawn err: %d",
        pData->needCleanUp, pData->spawnErr);
  if (!pData->needCleanUp || atomic_load_32(&pData->stopCalled)) {
    return 0;
  }
  atomic_store_32(&pData->stopCalled, 1);
  pData->needCleanUp = false;
  uv_barrier_destroy(&pData->barrier);
  uv_async_send(&pData->stopAsync);
  uv_thread_join(&pData->thread);
  fnInfo("dnode-mgmt udfd cleaned up");
  return 0;
}

//==============================================================================================
S
shenglian zhou 已提交
189 190 191 192
/* Copyright (c) 2013, Ben Noordhuis <info@bnoordhuis.nl>
 * The QUEUE is copied from queue.h under libuv
 * */

S
shenglian zhou 已提交
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 226 227 228 229 230 231 232 233 234 235 236 237 238 239 240 241 242 243 244 245 246 247 248 249 250 251 252 253 254 255 256 257 258 259 260 261 262 263 264 265 266 267 268 269 270 271 272 273 274 275 276 277 278 279 280
typedef void *QUEUE[2];

/* Private macros. */
#define QUEUE_NEXT(q)       (*(QUEUE **) &((*(q))[0]))
#define QUEUE_PREV(q)       (*(QUEUE **) &((*(q))[1]))
#define QUEUE_PREV_NEXT(q)  (QUEUE_NEXT(QUEUE_PREV(q)))
#define QUEUE_NEXT_PREV(q)  (QUEUE_PREV(QUEUE_NEXT(q)))

/* Public macros. */
#define QUEUE_DATA(ptr, type, field)                                          \
  ((type *) ((char *) (ptr) - offsetof(type, field)))

/* Important note: mutating the list while QUEUE_FOREACH is
 * iterating over its elements results in undefined behavior.
 */
#define QUEUE_FOREACH(q, h)                                                   \
  for ((q) = QUEUE_NEXT(h); (q) != (h); (q) = QUEUE_NEXT(q))

#define QUEUE_EMPTY(q)                                                        \
  ((const QUEUE *) (q) == (const QUEUE *) QUEUE_NEXT(q))

#define QUEUE_HEAD(q)                                                         \
  (QUEUE_NEXT(q))

#define QUEUE_INIT(q)                                                         \
  do {                                                                        \
    QUEUE_NEXT(q) = (q);                                                      \
    QUEUE_PREV(q) = (q);                                                      \
  }                                                                           \
  while (0)

#define QUEUE_ADD(h, n)                                                       \
  do {                                                                        \
    QUEUE_PREV_NEXT(h) = QUEUE_NEXT(n);                                       \
    QUEUE_NEXT_PREV(n) = QUEUE_PREV(h);                                       \
    QUEUE_PREV(h) = QUEUE_PREV(n);                                            \
    QUEUE_PREV_NEXT(h) = (h);                                                 \
  }                                                                           \
  while (0)

#define QUEUE_SPLIT(h, q, n)                                                  \
  do {                                                                        \
    QUEUE_PREV(n) = QUEUE_PREV(h);                                            \
    QUEUE_PREV_NEXT(n) = (n);                                                 \
    QUEUE_NEXT(n) = (q);                                                      \
    QUEUE_PREV(h) = QUEUE_PREV(q);                                            \
    QUEUE_PREV_NEXT(h) = (h);                                                 \
    QUEUE_PREV(q) = (n);                                                      \
  }                                                                           \
  while (0)

#define QUEUE_MOVE(h, n)                                                      \
  do {                                                                        \
    if (QUEUE_EMPTY(h))                                                       \
      QUEUE_INIT(n);                                                          \
    else {                                                                    \
      QUEUE* q = QUEUE_HEAD(h);                                               \
      QUEUE_SPLIT(h, q, n);                                                   \
    }                                                                         \
  }                                                                           \
  while (0)

#define QUEUE_INSERT_HEAD(h, q)                                               \
  do {                                                                        \
    QUEUE_NEXT(q) = QUEUE_NEXT(h);                                            \
    QUEUE_PREV(q) = (h);                                                      \
    QUEUE_NEXT_PREV(q) = (q);                                                 \
    QUEUE_NEXT(h) = (q);                                                      \
  }                                                                           \
  while (0)

#define QUEUE_INSERT_TAIL(h, q)                                               \
  do {                                                                        \
    QUEUE_NEXT(q) = (h);                                                      \
    QUEUE_PREV(q) = QUEUE_PREV(h);                                            \
    QUEUE_PREV_NEXT(q) = (q);                                                 \
    QUEUE_PREV(h) = (q);                                                      \
  }                                                                           \
  while (0)

#define QUEUE_REMOVE(q)                                                       \
  do {                                                                        \
    QUEUE_PREV_NEXT(q) = QUEUE_NEXT(q);                                       \
    QUEUE_NEXT_PREV(q) = QUEUE_PREV(q);                                       \
  }                                                                           \
  while (0)


281 282 283 284 285 286
enum {
  UV_TASK_CONNECT = 0,
  UV_TASK_REQ_RSP = 1,
  UV_TASK_DISCONNECT = 2
};

287 288
int64_t gUdfTaskSeqNum = 0;
typedef struct SUdfdProxy {
289
  char udfdPipeName[PATH_MAX + UDF_LISTEN_PIPE_NAME_LEN + 2];
290 291 292 293 294 295 296 297 298 299 300 301
  uv_barrier_t gUdfInitBarrier;

  uv_loop_t gUdfdLoop;
  uv_thread_t gUdfLoopThread;
  uv_async_t gUdfLoopTaskAync;

  uv_async_t gUdfLoopStopAsync;

  uv_mutex_t gUdfTaskQueueMutex;
  int8_t gUdfcState;
  QUEUE gUdfTaskQueue;
  QUEUE gUvProcTaskQueue;
302 303

  int8_t initialized;
304 305
} SUdfdProxy;

306
SUdfdProxy gUdfdProxy = {0};
307

S
slzhou 已提交
308
typedef struct SClientUdfUvSession {
309
  SUdfdProxy *udfc;
310
  int64_t severHandle;
S
slzhou 已提交
311
  uv_pipe_t *udfUvPipe;
S
shenglian zhou 已提交
312 313 314 315

  int8_t  outputType;
  int32_t outputLen;
  int32_t bufSize;
S
slzhou 已提交
316
} SClientUdfUvSession;
317 318

typedef struct SClientUvTaskNode {
319
  SUdfdProxy *udfc;
320 321 322 323 324 325 326 327 328 329 330
  int8_t type;
  int errCode;

  uv_pipe_t *pipe;

  int64_t seqNum;
  uv_buf_t reqBuf;

  uv_sem_t taskSem;
  uv_buf_t rspBuf;

S
shenglian zhou 已提交
331 332 333
  QUEUE recvTaskQueue;
  QUEUE procTaskQueue;
  QUEUE connTaskQueue;
334 335 336 337 338
} SClientUvTaskNode;

typedef struct SClientUdfTask {
  int8_t type;

S
slzhou 已提交
339
  SClientUdfUvSession *session;
340 341 342 343 344 345 346 347 348 349 350 351 352 353 354 355 356 357 358 359 360 361 362 363 364 365 366 367 368

  int32_t errCode;

  union {
    struct {
      SUdfSetupRequest req;
      SUdfSetupResponse rsp;
    } _setup;
    struct {
      SUdfCallRequest req;
      SUdfCallResponse rsp;
    } _call;
    struct {
      SUdfTeardownRequest req;
      SUdfTeardownResponse rsp;
    } _teardown;
  };

} SClientUdfTask;

typedef struct SClientConnBuf {
  char *buf;
  int32_t len;
  int32_t cap;
  int32_t total;
} SClientConnBuf;

typedef struct SClientUvConn {
  uv_pipe_t *pipe;
S
shenglian zhou 已提交
369
  QUEUE taskQueue;
370
  SClientConnBuf readBuf;
S
slzhou 已提交
371
  SClientUdfUvSession *session;
372 373
} SClientUvConn;

374 375
enum {
  UDFC_STATE_INITAL = 0, // initial state
376
  UDFC_STATE_STARTNG, // starting after udfcOpen
377
  UDFC_STATE_READY, // started and begin to receive quests
378
  UDFC_STATE_STOPPING, // stopping after udfcClose
379
};
380

381 382
int32_t getUdfdPipeName(char* pipeName, int32_t size) {
  char    dnodeId[8] = {0};
wafwerar's avatar
wafwerar 已提交
383
  size_t  dnodeIdSize = sizeof(dnodeId);
384 385
  int32_t err = uv_os_getenv(UDF_DNODE_ID_ENV_NAME, dnodeId, &dnodeIdSize);
  if (err != 0) {
386
    fnError("get dnode id from env. error: %s.", uv_err_name(err));
387 388
    dnodeId[0] = '1';
  }
389
#ifdef _WIN32
390
  snprintf(pipeName, size, "%s%s", UDF_LISTEN_PIPE_NAME_PREFIX, dnodeId);
391 392 393 394
#else
  snprintf(pipeName, size, "%s/%s%s", tsDataDir, UDF_LISTEN_PIPE_NAME_PREFIX, dnodeId);
#endif
  fnInfo("get dnode id from env. dnode id: %s. pipe path: %s", dnodeId, pipeName);
395 396 397
  return 0;
}

398 399 400
int32_t encodeUdfSetupRequest(void **buf, const SUdfSetupRequest *setup) {
  int32_t len = 0;
  len += taosEncodeBinary(buf, setup->udfName, TSDB_FUNC_NAME_LEN);
S
shenglian zhou 已提交
401 402
  return len;
}
403

404 405 406
void* decodeUdfSetupRequest(const void* buf, SUdfSetupRequest *request) {
  buf = taosDecodeBinaryTo(buf, request->udfName, TSDB_FUNC_NAME_LEN);
  return (void*)buf;
S
shenglian zhou 已提交
407
}
408

409 410
int32_t encodeUdfInterBuf(void **buf, const SUdfInterBuf* state) {
  int32_t len = 0;
411
  len += taosEncodeFixedI8(buf, state->numOfResult);
412 413
  len += taosEncodeFixedI32(buf, state->bufLen);
  len += taosEncodeBinary(buf, state->buf, state->bufLen);
S
shenglian zhou 已提交
414 415
  return len;
}
416

417
void* decodeUdfInterBuf(const void* buf, SUdfInterBuf* state) {
418
  buf = taosDecodeFixedI8(buf, &state->numOfResult);
419 420 421
  buf = taosDecodeFixedI32(buf, &state->bufLen);
  buf = taosDecodeBinary(buf, (void**)&state->buf, state->bufLen);
  return (void*)buf;
S
shenglian zhou 已提交
422 423
}

424 425 426 427 428 429 430 431 432 433 434 435 436 437 438 439
int32_t encodeUdfCallRequest(void **buf, const SUdfCallRequest *call) {
  int32_t len = 0;
  len += taosEncodeFixedI64(buf, call->udfHandle);
  len += taosEncodeFixedI8(buf, call->callType);
  if (call->callType == TSDB_UDF_CALL_SCALA_PROC) {
    len += tEncodeDataBlock(buf, &call->block);
  } else if (call->callType == TSDB_UDF_CALL_AGG_INIT) {
    len += taosEncodeFixedI8(buf, call->initFirst);
  } else if (call->callType == TSDB_UDF_CALL_AGG_PROC) {
    len += tEncodeDataBlock(buf, &call->block);
    len += encodeUdfInterBuf(buf, &call->interBuf);
  } else if (call->callType == TSDB_UDF_CALL_AGG_MERGE) {
    len += encodeUdfInterBuf(buf, &call->interBuf);
    len += encodeUdfInterBuf(buf, &call->interBuf2);
  } else if (call->callType == TSDB_UDF_CALL_AGG_FIN) {
    len += encodeUdfInterBuf(buf, &call->interBuf);
440
  }
441
  return len;
442 443
}

444 445 446 447 448 449 450 451 452 453 454 455 456 457 458 459 460 461 462 463 464
void* decodeUdfCallRequest(const void* buf, SUdfCallRequest* call) {
  buf = taosDecodeFixedI64(buf, &call->udfHandle);
  buf = taosDecodeFixedI8(buf, &call->callType);
  switch (call->callType) {
    case TSDB_UDF_CALL_SCALA_PROC:
      buf = tDecodeDataBlock(buf, &call->block);
      break;
    case TSDB_UDF_CALL_AGG_INIT:
      buf = taosDecodeFixedI8(buf, &call->initFirst);
      break;
    case TSDB_UDF_CALL_AGG_PROC:
      buf = tDecodeDataBlock(buf, &call->block);
      buf = decodeUdfInterBuf(buf, &call->interBuf);
      break;
    case TSDB_UDF_CALL_AGG_MERGE:
      buf = decodeUdfInterBuf(buf, &call->interBuf);
      buf = decodeUdfInterBuf(buf, &call->interBuf2);
      break;
    case TSDB_UDF_CALL_AGG_FIN:
      buf = decodeUdfInterBuf(buf, &call->interBuf);
      break;
465
  }
466
  return (void*)buf;
S
shenglian zhou 已提交
467 468
}

469 470 471 472
int32_t encodeUdfTeardownRequest(void **buf, const SUdfTeardownRequest *teardown) {
  int32_t len = 0;
  len += taosEncodeFixedI64(buf, teardown->udfHandle);
  return len;
S
shenglian zhou 已提交
473 474
}

475 476 477
void* decodeUdfTeardownRequest(const void* buf, SUdfTeardownRequest *teardown) {
  buf = taosDecodeFixedI64(buf, &teardown->udfHandle);
  return (void*)buf;
S
shenglian zhou 已提交
478 479
}

480 481 482 483 484 485 486 487 488 489 490 491 492 493 494 495 496 497
int32_t encodeUdfRequest(void** buf, const SUdfRequest* request) {
  int32_t len = 0;
  if (buf == NULL) {
    len += sizeof(request->msgLen);
  } else {
    *(int32_t*)(*buf) = request->msgLen;
    *buf = POINTER_SHIFT(*buf, sizeof(request->msgLen));
  }
  len += taosEncodeFixedI64(buf, request->seqNum);
  len += taosEncodeFixedI8(buf, request->type);
  if (request->type == UDF_TASK_SETUP) {
    len += encodeUdfSetupRequest(buf, &request->setup);
  } else if (request->type == UDF_TASK_CALL) {
    len += encodeUdfCallRequest(buf, &request->call);
  } else if (request->type == UDF_TASK_TEARDOWN) {
    len += encodeUdfTeardownRequest(buf, &request->teardown);
  }
  return len;
S
shenglian zhou 已提交
498 499
}

500 501
void* decodeUdfRequest(const void* buf, SUdfRequest* request) {
  request->msgLen = *(int32_t*)(buf);
S
slzhou 已提交
502
  buf = POINTER_SHIFT(buf, sizeof(request->msgLen));
S
shenglian zhou 已提交
503

504 505
  buf = taosDecodeFixedI64(buf, &request->seqNum);
  buf = taosDecodeFixedI8(buf, &request->type);
S
shenglian zhou 已提交
506 507

  if (request->type == UDF_TASK_SETUP) {
508
    buf = decodeUdfSetupRequest(buf, &request->setup);
S
shenglian zhou 已提交
509
  } else if (request->type == UDF_TASK_CALL) {
510 511 512
    buf = decodeUdfCallRequest(buf, &request->call);
  } else if (request->type == UDF_TASK_TEARDOWN) {
    buf = decodeUdfTeardownRequest(buf, &request->teardown);
S
shenglian zhou 已提交
513
  }
514
  return (void*)buf;
S
shenglian zhou 已提交
515
}
516

517 518 519
int32_t encodeUdfSetupResponse(void **buf, const SUdfSetupResponse *setupRsp) {
  int32_t len = 0;
  len += taosEncodeFixedI64(buf, setupRsp->udfHandle);
S
shenglian zhou 已提交
520 521 522
  len += taosEncodeFixedI8(buf, setupRsp->outputType);
  len += taosEncodeFixedI32(buf, setupRsp->outputLen);
  len += taosEncodeFixedI32(buf, setupRsp->bufSize);
523 524
  return len;
}
525

526 527
void* decodeUdfSetupResponse(const void* buf, SUdfSetupResponse* setupRsp) {
  buf = taosDecodeFixedI64(buf, &setupRsp->udfHandle);
S
shenglian zhou 已提交
528 529 530
  buf = taosDecodeFixedI8(buf, &setupRsp->outputType);
  buf = taosDecodeFixedI32(buf, &setupRsp->outputLen);
  buf = taosDecodeFixedI32(buf, &setupRsp->bufSize);
531
  return (void*)buf;
S
shenglian zhou 已提交
532
}
533

534 535 536 537 538 539
int32_t encodeUdfCallResponse(void **buf, const SUdfCallResponse *callRsp) {
  int32_t len = 0;
  len += taosEncodeFixedI8(buf, callRsp->callType);
  switch (callRsp->callType) {
    case TSDB_UDF_CALL_SCALA_PROC:
      len += tEncodeDataBlock(buf, &callRsp->resultData);
540
      break;
541
    case TSDB_UDF_CALL_AGG_INIT:
S
slzhou 已提交
542
      len += encodeUdfInterBuf(buf, &callRsp->resultBuf);
543 544
      break;
    case TSDB_UDF_CALL_AGG_PROC:
S
slzhou 已提交
545
      len += encodeUdfInterBuf(buf, &callRsp->resultBuf);
546 547
      break;
    case TSDB_UDF_CALL_AGG_MERGE:
S
slzhou 已提交
548
      len += encodeUdfInterBuf(buf, &callRsp->resultBuf);
549 550
      break;
    case TSDB_UDF_CALL_AGG_FIN:
S
slzhou 已提交
551
      len += encodeUdfInterBuf(buf, &callRsp->resultBuf);
552 553
      break;
  }
554
  return len;
S
shenglian zhou 已提交
555 556
}

557 558 559 560 561 562 563
void* decodeUdfCallResponse(const void* buf, SUdfCallResponse* callRsp) {
  buf = taosDecodeFixedI8(buf, &callRsp->callType);
  switch (callRsp->callType) {
    case TSDB_UDF_CALL_SCALA_PROC:
      buf = tDecodeDataBlock(buf, &callRsp->resultData);
      break;
    case TSDB_UDF_CALL_AGG_INIT:
S
slzhou 已提交
564
      buf = decodeUdfInterBuf(buf, &callRsp->resultBuf);
565 566
      break;
    case TSDB_UDF_CALL_AGG_PROC:
S
slzhou 已提交
567
      buf = decodeUdfInterBuf(buf, &callRsp->resultBuf);
568 569
      break;
    case TSDB_UDF_CALL_AGG_MERGE:
S
slzhou 已提交
570
      buf = decodeUdfInterBuf(buf, &callRsp->resultBuf);
571 572
      break;
    case TSDB_UDF_CALL_AGG_FIN:
S
slzhou 已提交
573
      buf = decodeUdfInterBuf(buf, &callRsp->resultBuf);
574
      break;
S
shenglian zhou 已提交
575
  }
576
  return (void*)buf;
S
shenglian zhou 已提交
577 578
}

579 580
int32_t encodeUdfTeardownResponse(void** buf, const SUdfTeardownResponse* teardownRsp) {
  return 0;
S
shenglian zhou 已提交
581 582
}

583 584
void* decodeUdfTeardownResponse(const void* buf, SUdfTeardownResponse* teardownResponse) {
  return (void*)buf;
S
shenglian zhou 已提交
585 586
}

587 588 589 590 591 592 593
int32_t encodeUdfResponse(void** buf, const SUdfResponse* rsp) {
  int32_t len = 0;
  if (buf == NULL) {
    len += sizeof(rsp->msgLen);
  } else {
    *(int32_t*)(*buf) = rsp->msgLen;
    *buf = POINTER_SHIFT(*buf, sizeof(rsp->msgLen));
S
shenglian zhou 已提交
594 595
  }

S
slzhou 已提交
596 597 598 599 600 601 602
  if (buf == NULL) {
    len += sizeof(rsp->seqNum);
  } else {
    *(int64_t*)(*buf) = rsp->seqNum;
    *buf = POINTER_SHIFT(*buf, sizeof(rsp->seqNum));
  }

603 604 605
  len += taosEncodeFixedI64(buf, rsp->seqNum);
  len += taosEncodeFixedI8(buf, rsp->type);
  len += taosEncodeFixedI32(buf, rsp->code);
S
shenglian zhou 已提交
606

607 608 609 610 611 612 613 614 615 616 617 618 619 620 621
  switch (rsp->type) {
    case UDF_TASK_SETUP:
      len += encodeUdfSetupResponse(buf, &rsp->setupRsp);
      break;
    case UDF_TASK_CALL:
      len += encodeUdfCallResponse(buf, &rsp->callRsp);
      break;
    case UDF_TASK_TEARDOWN:
      len += encodeUdfTeardownResponse(buf, &rsp->teardownRsp);
      break;
    default:
      //TODO: log error
      break;
  }
  return len;
S
shenglian zhou 已提交
622 623
}

624 625
void* decodeUdfResponse(const void* buf, SUdfResponse* rsp) {
  rsp->msgLen = *(int32_t*)(buf);
S
slzhou 已提交
626
  buf = POINTER_SHIFT(buf, sizeof(rsp->msgLen));
S
slzhou 已提交
627 628
  rsp->seqNum = *(int64_t*)(buf);
  buf = POINTER_SHIFT(buf, sizeof(rsp->seqNum));
629 630 631
  buf = taosDecodeFixedI64(buf, &rsp->seqNum);
  buf = taosDecodeFixedI8(buf, &rsp->type);
  buf = taosDecodeFixedI32(buf, &rsp->code);
S
shenglian zhou 已提交
632

633 634 635 636 637 638 639 640 641 642 643 644 645
  switch (rsp->type) {
    case UDF_TASK_SETUP:
      buf = decodeUdfSetupResponse(buf, &rsp->setupRsp);
      break;
    case UDF_TASK_CALL:
      buf = decodeUdfCallResponse(buf, &rsp->callRsp);
      break;
    case UDF_TASK_TEARDOWN:
      buf = decodeUdfTeardownResponse(buf, &rsp->teardownRsp);
      break;
    default:
      //TODO: log error
      break;
646
  }
647
  return (void*)buf;
648
}
649

S
shenglian zhou 已提交
650 651
void freeUdfColumnData(SUdfColumnData *data, SUdfColumnMeta *meta) {
  if (IS_VAR_DATA_TYPE(meta->type)) {
S
slzhou 已提交
652 653 654 655
    taosMemoryFree(data->varLenCol.varOffsets);
    data->varLenCol.varOffsets = NULL;
    taosMemoryFree(data->varLenCol.payload);
    data->varLenCol.payload = NULL;
S
shenglian zhou 已提交
656
  } else {
S
slzhou 已提交
657 658 659 660
    taosMemoryFree(data->fixLenCol.nullBitmap);
    data->fixLenCol.nullBitmap = NULL;
    taosMemoryFree(data->fixLenCol.data);
    data->fixLenCol.data = NULL;
S
shenglian zhou 已提交
661 662 663 664
  }
}

void freeUdfColumn(SUdfColumn* col) {
S
shenglian zhou 已提交
665
  freeUdfColumnData(&col->colData, &col->colMeta);
S
shenglian zhou 已提交
666 667 668 669 670 671 672 673 674 675 676 677 678 679 680 681 682
}

void freeUdfDataDataBlock(SUdfDataBlock *block) {
  for (int32_t i = 0; i < block->numOfCols; ++i) {
    freeUdfColumn(block->udfCols[i]);
    taosMemoryFree(block->udfCols[i]);
    block->udfCols[i] = NULL;
  }
  taosMemoryFree(block->udfCols);
  block->udfCols = NULL;
}

void freeUdfInterBuf(SUdfInterBuf *buf) {
  taosMemoryFree(buf->buf);
  buf->buf = NULL;
}

S
slzhou 已提交
683 684 685 686

int32_t convertDataBlockToUdfDataBlock(SSDataBlock *block, SUdfDataBlock *udfBlock) {
  udfBlock->numOfRows = block->info.rows;
  udfBlock->numOfCols = block->info.numOfCols;
S
slzhou 已提交
687
  udfBlock->udfCols = taosMemoryCalloc(udfBlock->numOfCols, sizeof(SUdfColumn*));
S
slzhou 已提交
688
  for (int32_t i = 0; i < udfBlock->numOfCols; ++i) {
S
slzhou 已提交
689
    udfBlock->udfCols[i] = taosMemoryCalloc(1, sizeof(SUdfColumn));
S
slzhou 已提交
690 691 692 693 694 695 696
    SColumnInfoData *col= (SColumnInfoData*)taosArrayGet(block->pDataBlock, i);
    SUdfColumn *udfCol = udfBlock->udfCols[i];
    udfCol->colMeta.type = col->info.type;
    udfCol->colMeta.bytes = col->info.bytes;
    udfCol->colMeta.scale = col->info.scale;
    udfCol->colMeta.precision = col->info.precision;
    udfCol->colData.numOfRows = udfBlock->numOfRows;
S
shenglian zhou 已提交
697
    if (IS_VAR_DATA_TYPE(udfCol->colMeta.type)) {
S
slzhou 已提交
698 699 700 701 702 703
      udfCol->colData.varLenCol.varOffsetsLen = sizeof(int32_t) * udfBlock->numOfRows;
      udfCol->colData.varLenCol.varOffsets = taosMemoryMalloc(udfCol->colData.varLenCol.varOffsetsLen);
      memcpy(udfCol->colData.varLenCol.varOffsets, col->varmeta.offset, udfCol->colData.varLenCol.varOffsetsLen);
      udfCol->colData.varLenCol.payloadLen = colDataGetLength(col, udfBlock->numOfRows);
      udfCol->colData.varLenCol.payload = taosMemoryMalloc(udfCol->colData.varLenCol.payloadLen);
      memcpy(udfCol->colData.varLenCol.payload, col->pData, udfCol->colData.varLenCol.payloadLen);
S
slzhou 已提交
704
    } else {
S
slzhou 已提交
705 706 707 708 709 710 711 712 713 714
      udfCol->colData.fixLenCol.nullBitmapLen = BitmapLen(udfCol->colData.numOfRows);
      int32_t bitmapLen = udfCol->colData.fixLenCol.nullBitmapLen;
      udfCol->colData.fixLenCol.nullBitmap = taosMemoryMalloc(udfCol->colData.fixLenCol.nullBitmapLen);
      char* bitmap = udfCol->colData.fixLenCol.nullBitmap;
      memcpy(bitmap, col->nullbitmap, bitmapLen);
      udfCol->colData.fixLenCol.dataLen = colDataGetLength(col, udfBlock->numOfRows);
      int32_t dataLen = udfCol->colData.fixLenCol.dataLen;
      udfCol->colData.fixLenCol.data = taosMemoryMalloc(udfCol->colData.fixLenCol.dataLen);
      char* data = udfCol->colData.fixLenCol.data;
      memcpy(data, col->pData, dataLen);
S
slzhou 已提交
715 716 717 718 719 720 721 722
    }
  }
  return 0;
}

int32_t convertUdfColumnToDataBlock(SUdfColumn *udfCol, SSDataBlock *block) {
  block->info.numOfCols = 1;
  block->info.rows = udfCol->colData.numOfRows;
S
shenglian zhou 已提交
723
  block->info.hasVarCol = IS_VAR_DATA_TYPE(udfCol->colMeta.type);
S
slzhou 已提交
724 725 726 727 728 729 730 731 732 733 734 735

  block->pDataBlock = taosArrayInit(1, sizeof(SColumnInfoData));
  taosArraySetSize(block->pDataBlock, 1);
  SColumnInfoData *col = taosArrayGet(block->pDataBlock, 0);
  SUdfColumnMeta *meta = &udfCol->colMeta;
  col->info.precision = meta->precision;
  col->info.bytes = meta->bytes;
  col->info.scale = meta->scale;
  col->info.type = meta->type;
  SUdfColumnData *data = &udfCol->colData;

  if (!IS_VAR_DATA_TYPE(meta->type)) {
S
slzhou 已提交
736 737 738 739
    col->nullbitmap = taosMemoryMalloc(data->fixLenCol.nullBitmapLen);
    memcpy(col->nullbitmap, data->fixLenCol.nullBitmap, data->fixLenCol.nullBitmapLen);
    col->pData = taosMemoryMalloc(data->fixLenCol.dataLen);
    memcpy(col->pData, data->fixLenCol.data, data->fixLenCol.dataLen);
S
slzhou 已提交
740
  } else {
S
slzhou 已提交
741 742 743 744
    col->varmeta.offset = taosMemoryMalloc(data->varLenCol.varOffsetsLen);
    memcpy(col->varmeta.offset, data->varLenCol.varOffsets, data->varLenCol.varOffsetsLen);
    col->pData = taosMemoryMalloc(data->varLenCol.payloadLen);
    memcpy(col->pData, data->varLenCol.payload, data->varLenCol.payloadLen);
S
slzhou 已提交
745 746 747 748
  }
  return 0;
}

S
slzhou 已提交
749 750 751 752 753 754 755 756 757 758 759 760 761 762
int32_t convertScalarParamToDataBlock(SScalarParam *input, int32_t numOfCols, SSDataBlock *output) {
  output->info.rows = input->numOfRows;
  output->info.numOfCols = numOfCols;
  bool hasVarCol = false;
  for (int32_t i = 0; i < numOfCols; ++i) {
    if (IS_VAR_DATA_TYPE((input+i)->columnData->info.type)) {
      hasVarCol = true;
      break;
    }
  }
  output->info.hasVarCol = hasVarCol;

  //TODO: free the array output->pDataBlock
  output->pDataBlock = taosArrayInit(numOfCols, sizeof(SColumnInfoData));
763 764 765
  for (int32_t i = 0; i < numOfCols; ++i) {
    taosArrayPush(output->pDataBlock, (input + i)->columnData);
  }
S
slzhou 已提交
766 767 768 769 770 771 772 773 774 775 776 777 778
  return 0;
}

int32_t convertDataBlockToScalarParm(SSDataBlock *input, SScalarParam *output) {
  if (input->info.numOfCols != 1) {
    fnError("scalar function only support one column");
    return -1;
  }
  output->numOfRows = input->info.rows;
  //TODO: memory
  output->columnData = taosArrayGet(input->pDataBlock, 0);
  return 0;
}
S
slzhou 已提交
779

780 781
void onUdfcPipeClose(uv_handle_t *handle) {
  SClientUvConn *conn = handle->data;
S
shenglian zhou 已提交
782 783 784
  if (!QUEUE_EMPTY(&conn->taskQueue)) {
    QUEUE* h = QUEUE_HEAD(&conn->taskQueue);
    SClientUvTaskNode *task = QUEUE_DATA(h, SClientUvTaskNode, connTaskQueue);
785
    task->errCode = 0;
S
shenglian zhou 已提交
786
    QUEUE_REMOVE(&task->procTaskQueue);
S
slzhou 已提交
787
    uv_sem_post(&task->taskSem);
788
  }
S
slzhou 已提交
789
  conn->session->udfUvPipe = NULL;
wafwerar's avatar
wafwerar 已提交
790 791 792
  taosMemoryFree(conn->readBuf.buf);
  taosMemoryFree(conn);
  taosMemoryFree((uv_pipe_t *) handle);
793 794
}

S
slzhou 已提交
795 796
int32_t udfcGetUdfTaskResultFromUvTask(SClientUdfTask *task, SClientUvTaskNode *uvTask) {
  fnDebug("udfc get uv task result. task: %p, uvTask: %p", task, uvTask);
797 798
  if (uvTask->type == UV_TASK_REQ_RSP) {
    if (uvTask->rspBuf.base != NULL) {
S
shenglian zhou 已提交
799
      SUdfResponse rsp;
800 801
      void* buf = decodeUdfResponse(uvTask->rspBuf.base, &rsp);
      assert(uvTask->rspBuf.len == POINTER_DISTANCE(buf, uvTask->rspBuf.base));
S
shenglian zhou 已提交
802
      task->errCode = rsp.code;
803 804 805

      switch (task->type) {
        case UDF_TASK_SETUP: {
S
shenglian zhou 已提交
806
          //TODO: copy or not
S
shenglian zhou 已提交
807
          task->_setup.rsp = rsp.setupRsp;
808 809 810
          break;
        }
        case UDF_TASK_CALL: {
S
shenglian zhou 已提交
811
          task->_call.rsp = rsp.callRsp;
S
shenglian zhou 已提交
812
          //TODO: copy or not
813 814 815
          break;
        }
        case UDF_TASK_TEARDOWN: {
S
shenglian zhou 已提交
816
          task->_teardown.rsp = rsp.teardownRsp;
817 818 819 820 821 822 823 824 825
          //TODO: copy or not?
          break;
        }
        default: {
          break;
        }
      }

      // TODO: the call buffer is setup and freed by udf invocation
wafwerar's avatar
wafwerar 已提交
826
      taosMemoryFree(uvTask->rspBuf.base);
827 828
    } else {
      task->errCode = uvTask->errCode;
829
    }
830 831 832 833
  } else if (uvTask->type == UV_TASK_CONNECT) {
    task->errCode = uvTask->errCode;
  } else if (uvTask->type == UV_TASK_DISCONNECT) {
    task->errCode = uvTask->errCode;
834
  }
835 836 837 838 839 840 841 842 843
  return 0;
}

void udfcAllocateBuffer(uv_handle_t *handle, size_t suggestedSize, uv_buf_t *buf) {
  SClientUvConn *conn = handle->data;
  SClientConnBuf *connBuf = &conn->readBuf;

  int32_t msgHeadSize = sizeof(int32_t) + sizeof(int64_t);
  if (connBuf->cap == 0) {
wafwerar's avatar
wafwerar 已提交
844
    connBuf->buf = taosMemoryMalloc(msgHeadSize);
845 846 847 848 849 850 851 852
    if (connBuf->buf) {
      connBuf->len = 0;
      connBuf->cap = msgHeadSize;
      connBuf->total = -1;

      buf->base = connBuf->buf;
      buf->len = connBuf->cap;
    } else {
853
      fnError("udfc allocate buffer failure. size: %d", msgHeadSize);
854 855 856 857 858
      buf->base = NULL;
      buf->len = 0;
    }
  } else {
    connBuf->cap = connBuf->total > connBuf->cap ? connBuf->total : connBuf->cap;
wafwerar's avatar
wafwerar 已提交
859
    void *resultBuf = taosMemoryRealloc(connBuf->buf, connBuf->cap);
860 861 862 863 864
    if (resultBuf) {
      connBuf->buf = resultBuf;
      buf->base = connBuf->buf + connBuf->len;
      buf->len = connBuf->cap - connBuf->len;
    } else {
865
      fnError("udfc re-allocate buffer failure. size: %d", connBuf->cap);
866 867 868 869 870
      buf->base = NULL;
      buf->len = 0;
    }
  }

871
  fnTrace("conn buf cap - len - total : %d - %d - %d", connBuf->cap, connBuf->len, connBuf->total);
872 873 874

}

875 876 877 878 879
bool isUdfcUvMsgComplete(SClientConnBuf *connBuf) {
  if (connBuf->total == -1 && connBuf->len >= sizeof(int32_t)) {
    connBuf->total = *(int32_t *) (connBuf->buf);
  }
  if (connBuf->len == connBuf->cap && connBuf->total == connBuf->cap) {
880
    fnTrace("udfc complete message is received, now handle it");
881
    return true;
882
  }
883 884 885 886 887
  return false;
}

void udfcUvHandleRsp(SClientUvConn *conn) {
  SClientConnBuf *connBuf = &conn->readBuf;
888
  int64_t seqNum = *(int64_t *) (connBuf->buf + sizeof(int32_t)); // msglen then seqnum
889

S
shenglian zhou 已提交
890
  if (QUEUE_EMPTY(&conn->taskQueue)) {
891
    fnError("udfc no task waiting for response on connection");
892 893 894 895
    return;
  }
  bool found = false;
  SClientUvTaskNode *taskFound = NULL;
S
shenglian zhou 已提交
896 897 898 899
  QUEUE* h = QUEUE_NEXT(&conn->taskQueue);
  SClientUvTaskNode *task = QUEUE_DATA(h, SClientUvTaskNode, connTaskQueue);

  while (h != &conn->taskQueue) {
900 901 902 903 904
    if (task->seqNum == seqNum) {
      if (found == false) {
        found = true;
        taskFound = task;
      } else {
905
        fnError("udfc more than one task waiting for the same response");
906 907
        continue;
      }
908
    }
S
shenglian zhou 已提交
909 910
    h = QUEUE_NEXT(h);
    task = QUEUE_DATA(h, SClientUvTaskNode, connTaskQueue);
911 912
  }

913 914
  if (taskFound) {
    taskFound->rspBuf = uv_buf_init(connBuf->buf, connBuf->len);
S
shenglian zhou 已提交
915 916
    QUEUE_REMOVE(&taskFound->connTaskQueue);
    QUEUE_REMOVE(&taskFound->procTaskQueue);
S
slzhou 已提交
917
    uv_sem_post(&taskFound->taskSem);
918
  } else {
919
    fnError("no task is waiting for the response.");
920 921 922 923 924 925
  }
  connBuf->buf = NULL;
  connBuf->total = -1;
  connBuf->len = 0;
  connBuf->cap = 0;
}
926

927
void udfcUvHandleError(SClientUvConn *conn) {
S
shenglian zhou 已提交
928 929 930 931
  while (!QUEUE_EMPTY(&conn->taskQueue)) {
    QUEUE* h = QUEUE_HEAD(&conn->taskQueue);
    SClientUvTaskNode *task = QUEUE_DATA(h, SClientUvTaskNode, connTaskQueue);
    task->errCode = UDFC_CODE_PIPE_READ_ERR;
S
slzhou 已提交
932
    QUEUE_REMOVE(&task->connTaskQueue);
S
shenglian zhou 已提交
933
    QUEUE_REMOVE(&task->procTaskQueue);
S
slzhou 已提交
934
    uv_sem_post(&task->taskSem);
S
shenglian zhou 已提交
935 936
  }

S
slzhou 已提交
937
  uv_close((uv_handle_t *) conn->pipe, onUdfcPipeClose);
938 939 940
}

void onUdfcRead(uv_stream_t *client, ssize_t nread, const uv_buf_t *buf) {
941
  fnTrace("udfc client %p, client read from pipe. nread: %zd", client, nread);
942 943 944 945 946 947 948 949 950 951 952 953
  if (nread == 0) return;

  SClientUvConn *conn = client->data;
  SClientConnBuf *connBuf = &conn->readBuf;
  if (nread > 0) {
    connBuf->len += nread;
    if (isUdfcUvMsgComplete(connBuf)) {
      udfcUvHandleRsp(conn);
    }

  }
  if (nread < 0) {
S
slzhou 已提交
954
    fnError("udfc client pipe %p read error: %zd, %s.", client, nread, uv_strerror(nread));
955
    if (nread == UV_EOF) {
S
slzhou 已提交
956
      fnError("\tudfc client pipe %p closed", client);
957 958
    }
    udfcUvHandleError(conn);
959 960 961 962
  }

}

963 964
void onUdfClientWrite(uv_write_t *write, int status) {
  SClientUvTaskNode *uvTask = write->data;
965
  uv_pipe_t *pipe = uvTask->pipe;
966 967
  if (status == 0) {
    SClientUvConn *conn = pipe->data;
S
shenglian zhou 已提交
968
    QUEUE_INSERT_TAIL(&conn->taskQueue, &uvTask->connTaskQueue);
969
  } else {
970
    fnError("udfc client %p write error.", pipe);
971
  }
972
  fnTrace("udfc client %p write length:%zu", pipe, uvTask->reqBuf.len);
wafwerar's avatar
wafwerar 已提交
973 974
  taosMemoryFree(write);
  taosMemoryFree(uvTask->reqBuf.base);
975
}
H
Haojun Liao 已提交
976

977 978 979 980 981
void onUdfClientConnect(uv_connect_t *connect, int status) {
  SClientUvTaskNode *uvTask = connect->data;
  uvTask->errCode = status;
  if (status != 0) {
    //TODO: LOG error
H
Haojun Liao 已提交
982
  }
983
  uv_read_start((uv_stream_t *) uvTask->pipe, udfcAllocateBuffer, onUdfcRead);
wafwerar's avatar
wafwerar 已提交
984
  taosMemoryFree(connect);
985
  uv_sem_post(&uvTask->taskSem);
S
shenglian zhou 已提交
986
  QUEUE_REMOVE(&uvTask->procTaskQueue);
987
}
H
Haojun Liao 已提交
988

S
slzhou 已提交
989
int32_t udfcCreateUvTask(SClientUdfTask *task, int8_t uvTaskType, SClientUvTaskNode **pUvTask) {
wafwerar's avatar
wafwerar 已提交
990
  SClientUvTaskNode *uvTask = taosMemoryCalloc(1, sizeof(SClientUvTaskNode));
991
  uvTask->type = uvTaskType;
992
  uvTask->udfc = task->session->udfc;
993 994 995

  if (uvTaskType == UV_TASK_CONNECT) {
  } else if (uvTaskType == UV_TASK_REQ_RSP) {
S
slzhou 已提交
996
    uvTask->pipe = task->session->udfUvPipe;
997 998
    SUdfRequest request;
    request.type = task->type;
999
    request.seqNum =   atomic_fetch_add_64(&gUdfTaskSeqNum, 1);
1000 1001

    if (task->type == UDF_TASK_SETUP) {
S
shenglian zhou 已提交
1002
      request.setup = task->_setup.req;
1003 1004
      request.type = UDF_TASK_SETUP;
    } else if (task->type == UDF_TASK_CALL) {
S
shenglian zhou 已提交
1005
      request.call = task->_call.req;
1006 1007
      request.type = UDF_TASK_CALL;
    } else if (task->type == UDF_TASK_TEARDOWN) {
S
shenglian zhou 已提交
1008
      request.teardown = task->_teardown.req;
1009 1010 1011 1012
      request.type = UDF_TASK_TEARDOWN;
    } else {
      //TODO log and return error
    }
1013 1014
    int32_t bufLen = encodeUdfRequest(NULL, &request);
    request.msgLen = bufLen;
S
slzhou 已提交
1015 1016
    void *bufBegin = taosMemoryMalloc(bufLen);
    void *buf = bufBegin;
1017
    encodeUdfRequest(&buf, &request);
S
slzhou 已提交
1018
    uvTask->reqBuf = uv_buf_init(bufBegin, bufLen);
1019 1020
    uvTask->seqNum = request.seqNum;
  } else if (uvTaskType == UV_TASK_DISCONNECT) {
S
slzhou 已提交
1021
    uvTask->pipe = task->session->udfUvPipe;
1022 1023
  }
  uv_sem_init(&uvTask->taskSem, 0);
H
Haojun Liao 已提交
1024

1025 1026 1027
  *pUvTask = uvTask;
  return 0;
}
H
Haojun Liao 已提交
1028

S
slzhou 已提交
1029
int32_t udfcQueueUvTask(SClientUvTaskNode *uvTask) {
1030
  fnTrace("queue uv task to event loop, task: %d, %p", uvTask->type, uvTask);
1031 1032 1033 1034 1035
  SUdfdProxy *udfc = uvTask->udfc;
  uv_mutex_lock(&udfc->gUdfTaskQueueMutex);
  QUEUE_INSERT_TAIL(&udfc->gUdfTaskQueue, &uvTask->recvTaskQueue);
  uv_mutex_unlock(&udfc->gUdfTaskQueueMutex);
  uv_async_send(&udfc->gUdfLoopTaskAync);
H
Haojun Liao 已提交
1036

1037
  uv_sem_wait(&uvTask->taskSem);
S
slzhou 已提交
1038
  fnInfo("udfc uv task finished. task: %d, %p", uvTask->type, uvTask);
1039
  uv_sem_destroy(&uvTask->taskSem);
H
Haojun Liao 已提交
1040

1041 1042
  return 0;
}
H
Haojun Liao 已提交
1043

S
slzhou 已提交
1044
int32_t udfcStartUvTask(SClientUvTaskNode *uvTask) {
1045
  fnTrace("event loop start uv task. task: %d, %p", uvTask->type, uvTask);
1046 1047
  switch (uvTask->type) {
    case UV_TASK_CONNECT: {
wafwerar's avatar
wafwerar 已提交
1048
      uv_pipe_t *pipe = taosMemoryMalloc(sizeof(uv_pipe_t));
1049
      uv_pipe_init(&uvTask->udfc->gUdfdLoop, pipe, 0);
1050
      uvTask->pipe = pipe;
H
Haojun Liao 已提交
1051

S
slzhou 已提交
1052
      SClientUvConn *conn = taosMemoryCalloc(1, sizeof(SClientUvConn));
1053 1054 1055 1056 1057
      conn->pipe = pipe;
      conn->readBuf.len = 0;
      conn->readBuf.cap = 0;
      conn->readBuf.buf = 0;
      conn->readBuf.total = -1;
S
shenglian zhou 已提交
1058
      QUEUE_INIT(&conn->taskQueue);
H
Haojun Liao 已提交
1059

1060 1061
      pipe->data = conn;

wafwerar's avatar
wafwerar 已提交
1062
      uv_connect_t *connReq = taosMemoryMalloc(sizeof(uv_connect_t));
1063
      connReq->data = uvTask;
1064
      uv_pipe_connect(connReq, pipe, uvTask->udfc->udfdPipeName, onUdfClientConnect);
H
Haojun Liao 已提交
1065
      break;
1066 1067 1068
    }
    case UV_TASK_REQ_RSP: {
      uv_pipe_t *pipe = uvTask->pipe;
wafwerar's avatar
wafwerar 已提交
1069
      uv_write_t *write = taosMemoryMalloc(sizeof(uv_write_t));
1070 1071 1072 1073 1074 1075
      write->data = uvTask;
      uv_write(write, (uv_stream_t *) pipe, &uvTask->reqBuf, 1, onUdfClientWrite);
      break;
    }
    case UV_TASK_DISCONNECT: {
      SClientUvConn *conn = uvTask->pipe->data;
S
shenglian zhou 已提交
1076
      QUEUE_INSERT_TAIL(&conn->taskQueue, &uvTask->connTaskQueue);
1077 1078 1079 1080 1081 1082 1083
      uv_close((uv_handle_t *) uvTask->pipe, onUdfcPipeClose);
      break;
    }
    default: {
      break;
    }
  }
H
Haojun Liao 已提交
1084

1085 1086
  return 0;
}
H
Haojun Liao 已提交
1087

1088
void udfClientAsyncCb(uv_async_t *async) {
1089
  SUdfdProxy *udfc = async->data;
S
shenglian zhou 已提交
1090
  QUEUE wq;
1091

1092 1093 1094
  uv_mutex_lock(&udfc->gUdfTaskQueueMutex);
  QUEUE_MOVE(&udfc->gUdfTaskQueue, &wq);
  uv_mutex_unlock(&udfc->gUdfTaskQueueMutex);
1095

S
shenglian zhou 已提交
1096 1097 1098 1099
  while (!QUEUE_EMPTY(&wq)) {
    QUEUE* h = QUEUE_HEAD(&wq);
    QUEUE_REMOVE(h);
    SClientUvTaskNode *task = QUEUE_DATA(h, SClientUvTaskNode, recvTaskQueue);
S
slzhou 已提交
1100
    udfcStartUvTask(task);
1101
    QUEUE_INSERT_TAIL(&udfc->gUvProcTaskQueue, &task->procTaskQueue);
1102 1103 1104 1105
  }

}

1106
void cleanUpUvTasks(SUdfdProxy *udfc) {
S
slzhou 已提交
1107
  fnDebug("clean up uv tasks")
S
shenglian zhou 已提交
1108
  QUEUE wq;
1109

1110 1111 1112
  uv_mutex_lock(&udfc->gUdfTaskQueueMutex);
  QUEUE_MOVE(&udfc->gUdfTaskQueue, &wq);
  uv_mutex_unlock(&udfc->gUdfTaskQueueMutex);
1113

S
shenglian zhou 已提交
1114 1115 1116 1117
  while (!QUEUE_EMPTY(&wq)) {
    QUEUE* h = QUEUE_HEAD(&wq);
    QUEUE_REMOVE(h);
    SClientUvTaskNode *task = QUEUE_DATA(h, SClientUvTaskNode, recvTaskQueue);
1118
    if (udfc->gUdfcState == UDFC_STATE_STOPPING) {
1119 1120 1121 1122 1123
      task->errCode = UDFC_CODE_STOPPING;
    }
    uv_sem_post(&task->taskSem);
  }

1124 1125
  while (!QUEUE_EMPTY(&udfc->gUvProcTaskQueue)) {
    QUEUE* h = QUEUE_HEAD(&udfc->gUvProcTaskQueue);
S
shenglian zhou 已提交
1126 1127
    QUEUE_REMOVE(h);
    SClientUvTaskNode *task = QUEUE_DATA(h, SClientUvTaskNode, procTaskQueue);
1128
    if (udfc->gUdfcState == UDFC_STATE_STOPPING) {
S
shenglian zhou 已提交
1129 1130 1131 1132 1133
      task->errCode = UDFC_CODE_STOPPING;
    }
    uv_sem_post(&task->taskSem);
  }
}
1134

S
shenglian zhou 已提交
1135
void udfStopAsyncCb(uv_async_t *async) {
1136 1137 1138 1139
  SUdfdProxy *udfc = async->data;
  cleanUpUvTasks(udfc);
  if (udfc->gUdfcState == UDFC_STATE_STOPPING) {
    uv_stop(&udfc->gUdfdLoop);
S
shenglian zhou 已提交
1140
  }
1141
}
S
shenglian zhou 已提交
1142

S
shenglian zhou 已提交
1143
void constructUdfService(void *argsThread) {
1144 1145 1146 1147 1148 1149 1150 1151 1152 1153 1154
  SUdfdProxy *udfc = (SUdfdProxy*)argsThread;
  uv_loop_init(&udfc->gUdfdLoop);

  uv_async_init(&udfc->gUdfdLoop, &udfc->gUdfLoopTaskAync, udfClientAsyncCb);
  udfc->gUdfLoopTaskAync.data = udfc;
  uv_async_init(&udfc->gUdfdLoop, &udfc->gUdfLoopStopAsync, udfStopAsyncCb);
  udfc->gUdfLoopStopAsync.data = udfc;
  uv_mutex_init(&udfc->gUdfTaskQueueMutex);
  QUEUE_INIT(&udfc->gUdfTaskQueue);
  QUEUE_INIT(&udfc->gUvProcTaskQueue);
  uv_barrier_wait(&udfc->gUdfInitBarrier);
1155
  //TODO return value of uv_run
1156 1157
  uv_run(&udfc->gUdfdLoop, UV_RUN_DEFAULT);
  uv_loop_close(&udfc->gUdfdLoop);
1158 1159
}

1160 1161 1162 1163 1164 1165
int32_t udfcOpen() {
  int8_t old = atomic_val_compare_exchange_8(&gUdfdProxy.initialized, 0, 1);
  if (old == 1) {
    return 0;
  }
  SUdfdProxy *proxy = &gUdfdProxy;
1166
  getUdfdPipeName(proxy->udfdPipeName, sizeof(proxy->udfdPipeName));
1167 1168 1169
  proxy->gUdfcState = UDFC_STATE_STARTNG;
  uv_barrier_init(&proxy->gUdfInitBarrier, 2);
  uv_thread_create(&proxy->gUdfLoopThread, constructUdfService, proxy);
1170
  atomic_store_8(&proxy->gUdfcState, UDFC_STATE_READY);
1171
  proxy->gUdfcState = UDFC_STATE_READY;
1172 1173
  uv_barrier_wait(&proxy->gUdfInitBarrier);
  fnInfo("udfc initialized")
1174 1175 1176
  return 0;
}

1177 1178 1179 1180 1181 1182 1183
int32_t udfcClose() {
  int8_t old = atomic_val_compare_exchange_8(&gUdfdProxy.initialized, 1, 0);
  if (old == 0) {
    return 0;
  }

  SUdfdProxy *udfc = &gUdfdProxy;
1184 1185 1186 1187 1188
  udfc->gUdfcState = UDFC_STATE_STOPPING;
  uv_async_send(&udfc->gUdfLoopStopAsync);
  uv_thread_join(&udfc->gUdfLoopThread);
  uv_mutex_destroy(&udfc->gUdfTaskQueueMutex);
  uv_barrier_destroy(&udfc->gUdfInitBarrier);
1189 1190
  udfc->gUdfcState = UDFC_STATE_INITAL;
  fnInfo("udfc cleaned up");
1191 1192 1193
  return 0;
}

S
slzhou 已提交
1194
int32_t udfcRunUdfUvTask(SClientUdfTask *task, int8_t uvTaskType) {
1195 1196
  SClientUvTaskNode *uvTask = NULL;

S
slzhou 已提交
1197 1198 1199
  udfcCreateUvTask(task, uvTaskType, &uvTask);
  udfcQueueUvTask(uvTask);
  udfcGetUdfTaskResultFromUvTask(task, uvTask);
1200
  if (uvTaskType == UV_TASK_CONNECT) {
S
slzhou 已提交
1201 1202 1203
    task->session->udfUvPipe = uvTask->pipe;
    SClientUvConn *conn = uvTask->pipe->data;
    conn->session = task->session;
S
slzhou 已提交
1204 1205
  }
  taosMemoryFree(uvTask);
1206 1207 1208 1209
  uvTask = NULL;
  return task->errCode;
}

1210 1211 1212 1213 1214
int32_t setupUdf(char udfName[], UdfcFuncHandle *funcHandle) {
  fnInfo("udfc setup udf. udfName: %s", udfName);
  if (gUdfdProxy.gUdfcState != UDFC_STATE_READY) {
    return UDFC_CODE_INVALID_STATE;
  }
S
slzhou 已提交
1215
  SClientUdfTask *task = taosMemoryCalloc(1,sizeof(SClientUdfTask));
1216
  task->errCode = 0;
S
slzhou 已提交
1217
  task->session = taosMemoryCalloc(1, sizeof(SClientUdfUvSession));
1218
  task->session->udfc = &gUdfdProxy;
1219 1220 1221
  task->type = UDF_TASK_SETUP;

  SUdfSetupRequest *req = &task->_setup.req;
S
shenglian zhou 已提交
1222
  memcpy(req->udfName, udfName, TSDB_FUNC_NAME_LEN);
1223

S
slzhou 已提交
1224
  int32_t errCode = udfcRunUdfUvTask(task, UV_TASK_CONNECT);
1225
  if (errCode != 0) {
1226 1227
    fnError("failed to connect to pipe. udfName: %s, pipe: %s", udfName, (&gUdfdProxy)->udfdPipeName);
    return UDFC_CODE_CONNECT_PIPE_ERR;
H
Haojun Liao 已提交
1228
  }
1229

S
slzhou 已提交
1230
  udfcRunUdfUvTask(task, UV_TASK_REQ_RSP);
1231 1232 1233

  SUdfSetupResponse *rsp = &task->_setup.rsp;
  task->session->severHandle = rsp->udfHandle;
S
shenglian zhou 已提交
1234 1235 1236
  task->session->outputType = rsp->outputType;
  task->session->outputLen = rsp->outputLen;
  task->session->bufSize = rsp->bufSize;
1237 1238 1239 1240 1241 1242
  if (task->errCode != 0) {
    fnError("failed to setup udf. err: %d", task->errCode)
  } else {
    fnInfo("sucessfully setup udf func handle. handle: %p", task->session);
    *funcHandle = task->session;
  }
1243
  int32_t err = task->errCode;
wafwerar's avatar
wafwerar 已提交
1244
  taosMemoryFree(task);
1245
  return err;
H
Haojun Liao 已提交
1246 1247
}

1248
int32_t callUdf(UdfcFuncHandle handle, int8_t callType, SSDataBlock *input, SUdfInterBuf *state, SUdfInterBuf *state2,
1249
                SSDataBlock* output, SUdfInterBuf *newState) {
1250
  fnTrace("udfc call udf. callType: %d, funcHandle: %p", callType, handle);
S
slzhou 已提交
1251 1252 1253 1254 1255 1256
  SClientUdfUvSession *session = (SClientUdfUvSession *) handle;
  if (session->udfUvPipe == NULL) {
    fnError("No pipe to udfd");
    return UDFC_CODE_NO_PIPE;
  }
  SClientUdfTask *task = taosMemoryCalloc(1, sizeof(SClientUdfTask));
1257
  task->errCode = 0;
S
slzhou 已提交
1258
  task->session = (SClientUdfUvSession *) handle;
1259 1260 1261
  task->type = UDF_TASK_CALL;

  SUdfCallRequest *req = &task->_call.req;
S
slzhou 已提交
1262
  req->udfHandle = task->session->severHandle;
S
slzhou 已提交
1263
  req->callType = callType;
S
slzhou 已提交
1264

S
shenglian zhou 已提交
1265
  switch (callType) {
1266 1267 1268 1269
    case TSDB_UDF_CALL_AGG_INIT: {
      req->initFirst = 1;
      break;
    }
S
shenglian zhou 已提交
1270 1271 1272 1273 1274
    case TSDB_UDF_CALL_AGG_PROC: {
      req->block = *input;
      req->interBuf = *state;
      break;
    }
1275 1276 1277 1278 1279 1280
    case TSDB_UDF_CALL_AGG_MERGE: {
      req->interBuf = *state;
      req->interBuf2 = *state2;
      break;
    }
    case TSDB_UDF_CALL_AGG_FIN: {
S
shenglian zhou 已提交
1281 1282 1283 1284 1285 1286 1287 1288 1289
      req->interBuf = *state;
      break;
    }
    case TSDB_UDF_CALL_SCALA_PROC: {
      req->block = *input;
      break;
    }
  }

S
slzhou 已提交
1290
  udfcRunUdfUvTask(task, UV_TASK_REQ_RSP);
1291

1292 1293 1294 1295 1296 1297 1298 1299 1300 1301 1302 1303 1304 1305 1306 1307 1308 1309 1310 1311 1312 1313 1314 1315 1316
  if (task->errCode != 0) {
    fnError("call udf failure. err: %d", task->errCode);
  } else {
    SUdfCallResponse *rsp = &task->_call.rsp;
    switch (callType) {
      case TSDB_UDF_CALL_AGG_INIT: {
        *newState = rsp->resultBuf;
        break;
      }
      case TSDB_UDF_CALL_AGG_PROC: {
        *newState = rsp->resultBuf;
        break;
      }
      case TSDB_UDF_CALL_AGG_MERGE: {
        *newState = rsp->resultBuf;
        break;
      }
      case TSDB_UDF_CALL_AGG_FIN: {
        *newState = rsp->resultBuf;
        break;
      }
      case TSDB_UDF_CALL_SCALA_PROC: {
        *output = rsp->resultData;
        break;
      }
S
shenglian zhou 已提交
1317
    }
S
slzhou 已提交
1318 1319
  };
  int err = task->errCode;
wafwerar's avatar
wafwerar 已提交
1320
  taosMemoryFree(task);
S
slzhou 已提交
1321
  return err;
1322 1323
}

1324
int32_t callUdfAggInit(UdfcFuncHandle handle, SUdfInterBuf *interBuf) {
S
slzhou 已提交
1325 1326 1327 1328 1329 1330 1331 1332 1333
  int8_t callType = TSDB_UDF_CALL_AGG_INIT;

  int32_t err = callUdf(handle, callType, NULL, NULL, NULL, NULL, interBuf);

  return err;
}

// input: block, state
// output: interbuf,
1334
int32_t callUdfAggProcess(UdfcFuncHandle handle, SSDataBlock *block, SUdfInterBuf *state, SUdfInterBuf *newState) {
S
slzhou 已提交
1335 1336 1337 1338 1339 1340 1341
  int8_t callType = TSDB_UDF_CALL_AGG_PROC;
  int32_t err = callUdf(handle, callType, block, state, NULL, NULL, newState);
  return err;
}

// input: interbuf1, interbuf2
// output: resultBuf
1342
int32_t callUdfAggMerge(UdfcFuncHandle handle, SUdfInterBuf *interBuf1, SUdfInterBuf *interBuf2, SUdfInterBuf *resultBuf) {
S
slzhou 已提交
1343 1344 1345 1346 1347 1348 1349
  int8_t callType = TSDB_UDF_CALL_AGG_MERGE;
  int32_t err = callUdf(handle, callType, NULL, interBuf1, interBuf2, NULL, resultBuf);
  return err;
}

// input: interBuf
// output: resultData
1350
int32_t callUdfAggFinalize(UdfcFuncHandle handle, SUdfInterBuf *interBuf, SUdfInterBuf *resultData) {
S
slzhou 已提交
1351
  int8_t callType = TSDB_UDF_CALL_AGG_FIN;
S
slzhou 已提交
1352 1353 1354 1355
  int32_t err = callUdf(handle, callType, NULL, interBuf, NULL, NULL, resultData);
  return err;
}

S
slzhou 已提交
1356
int32_t callUdfScalarFunc(UdfcFuncHandle handle, SScalarParam *input, int32_t numOfCols, SScalarParam* output) {
S
slzhou 已提交
1357
  int8_t callType = TSDB_UDF_CALL_SCALA_PROC;
S
slzhou 已提交
1358 1359 1360 1361
  SSDataBlock inputBlock = {0};
  convertScalarParamToDataBlock(input, numOfCols, &inputBlock);
  SSDataBlock resultBlock = {0};
  int32_t err = callUdf(handle, callType, &inputBlock, NULL, NULL, &resultBlock, NULL);
S
slzhou 已提交
1362 1363 1364
  if (err == 0) {
    convertDataBlockToScalarParm(&resultBlock, output);
  }
S
slzhou 已提交
1365 1366 1367
  return err;
}

1368
int32_t teardownUdf(UdfcFuncHandle handle) {
1369
  fnInfo("tear down udf. udf func handle: %p", handle);
1370

S
slzhou 已提交
1371 1372 1373 1374 1375 1376 1377
  SClientUdfUvSession *session = (SClientUdfUvSession *) handle;
  if (session->udfUvPipe == NULL) {
    fnError("pipe to udfd does not exist");
    return UDFC_CODE_NO_PIPE;
  }

  SClientUdfTask *task = taosMemoryCalloc(1, sizeof(SClientUdfTask));
1378
  task->errCode = 0;
S
slzhou 已提交
1379
  task->session = session;
1380 1381 1382 1383 1384
  task->type = UDF_TASK_TEARDOWN;

  SUdfTeardownRequest *req = &task->_teardown.req;
  req->udfHandle = task->session->severHandle;

S
slzhou 已提交
1385
  udfcRunUdfUvTask(task, UV_TASK_REQ_RSP);
1386 1387 1388 1389 1390

  SUdfTeardownResponse *rsp = &task->_teardown.rsp;

  int32_t err = task->errCode;

S
slzhou 已提交
1391
  udfcRunUdfUvTask(task, UV_TASK_DISCONNECT);
1392

wafwerar's avatar
wafwerar 已提交
1393 1394
  taosMemoryFree(task->session);
  taosMemoryFree(task);
1395 1396 1397

  return err;
}
S
shenglian zhou 已提交
1398

S
shenglian zhou 已提交
1399 1400
//memory layout |---SUdfAggRes----|-----final result-----|---inter result----|
typedef struct SUdfAggRes {
S
slzhou 已提交
1401
  SClientUdfUvSession *session;
S
shenglian zhou 已提交
1402 1403 1404 1405 1406 1407
  int8_t finalResNum;
  int8_t interResNum;
  char* finalResBuf;
  char* interResBuf;
} SUdfAggRes;

S
shenglian zhou 已提交
1408
bool udfAggGetEnv(struct SFunctionNode* pFunc, SFuncExecEnv* pEnv) {
S
slzhou 已提交
1409
  if (fmIsScalarFunc(pFunc->funcId)) {
S
shenglian zhou 已提交
1410 1411
    return false;
  }
S
slzhou 已提交
1412
  pEnv->calcMemSize = sizeof(SUdfAggRes) + pFunc->node.resType.bytes + pFunc->udfBufSize;
S
shenglian zhou 已提交
1413 1414 1415 1416 1417 1418 1419 1420
  return true;
}

bool udfAggInit(struct SqlFunctionCtx *pCtx, struct SResultRowEntryInfo* pResultCellInfo) {
  if (functionSetup(pCtx, pResultCellInfo) != true) {
    return false;
  }
  UdfcFuncHandle handle;
1421 1422 1423
  int32_t udfCode = 0;
  if ((udfCode = setupUdf((char*)pCtx->udfName, &handle)) != 0) {
    fnError("udfAggInit error. step setupUdf. udf code: %d", udfCode);
S
shenglian zhou 已提交
1424 1425
    return false;
  }
S
slzhou 已提交
1426
  SClientUdfUvSession *session = (SClientUdfUvSession *)handle;
S
shenglian zhou 已提交
1427 1428
  SUdfAggRes *udfRes = (SUdfAggRes*)GET_ROWCELL_INTERBUF(pResultCellInfo);
  int32_t envSize = sizeof(SUdfAggRes) + session->outputLen + session->bufSize;
S
shenglian zhou 已提交
1429
  memset(udfRes, 0, envSize);
S
shenglian zhou 已提交
1430

S
slzhou 已提交
1431 1432 1433
  udfRes->finalResBuf = (char*)udfRes + sizeof(SUdfAggRes);
  udfRes->interResBuf = (char*)udfRes + sizeof(SUdfAggRes) + session->outputLen;

S
slzhou 已提交
1434
  udfRes->session = (SClientUdfUvSession *)handle;
S
shenglian zhou 已提交
1435
  SUdfInterBuf buf = {0};
1436 1437
  if ((udfCode = callUdfAggInit(handle, &buf)) != 0) {
    fnError("udfAggInit error. step callUdfAggInit. udf code: %d", udfCode);
S
shenglian zhou 已提交
1438 1439
    return false;
  }
S
shenglian zhou 已提交
1440
  udfRes->interResNum = buf.numOfResult;
S
slzhou 已提交
1441
  memcpy(udfRes->interResBuf, buf.buf, buf.bufLen);
S
shenglian zhou 已提交
1442 1443 1444 1445 1446 1447 1448
  return true;
}

int32_t udfAggProcess(struct SqlFunctionCtx *pCtx) {
  SInputColumnInfoData* pInput = &pCtx->input;
  int32_t numOfCols = pInput->numOfInputCols;

S
slzhou 已提交
1449
  SUdfAggRes* udfRes = (SUdfAggRes *)GET_ROWCELL_INTERBUF(GET_RES_INFO(pCtx));
S
slzhou 已提交
1450
  SClientUdfUvSession *session = udfRes->session;
S
slzhou 已提交
1451 1452
  udfRes->finalResBuf = (char*)udfRes + sizeof(SUdfAggRes);
  udfRes->interResBuf = (char*)udfRes + sizeof(SUdfAggRes) + session->outputLen;
S
shenglian zhou 已提交
1453 1454 1455 1456 1457

  int32_t start = pInput->startRowIndex;
  int32_t numOfRows = pInput->numOfRows;


1458 1459 1460 1461
  SSDataBlock tempBlock = {0};
  tempBlock.info.numOfCols = numOfCols;
  tempBlock.info.rows = numOfRows;
  tempBlock.info.uid = pInput->uid;
S
shenglian zhou 已提交
1462
  bool hasVarCol = false;
1463
  tempBlock.pDataBlock = taosArrayInit(numOfCols, sizeof(SColumnInfoData));
S
shenglian zhou 已提交
1464 1465 1466 1467 1468 1469

  for (int32_t i = 0; i < numOfCols; ++i) {
    SColumnInfoData *col = pInput->pData[i];
    if (IS_VAR_DATA_TYPE(col->info.type)) {
      hasVarCol = true;
    }
1470
    taosArrayPush(tempBlock.pDataBlock, col);
S
shenglian zhou 已提交
1471
  }
1472
  tempBlock.info.hasVarCol = hasVarCol;
S
shenglian zhou 已提交
1473

1474 1475
  SSDataBlock *inputBlock = blockDataExtractBlock(&tempBlock, start, numOfRows);

S
slzhou 已提交
1476
  SUdfInterBuf state = {.buf = udfRes->interResBuf,
1477
                        .bufLen = session->bufSize,
S
slzhou 已提交
1478
                        .numOfResult = udfRes->interResNum};
S
shenglian zhou 已提交
1479
  SUdfInterBuf newState = {0};
1480

1481 1482 1483 1484 1485 1486 1487 1488
  int32_t udfCode = callUdfAggProcess(session, inputBlock, &state, &newState);
  if (udfCode != 0) {
    fnError("udfAggProcess error. code: %d", udfCode);
    newState.numOfResult = 0;
  } else {
    udfRes->interResNum = newState.numOfResult;
    memcpy(udfRes->interResBuf, newState.buf, newState.bufLen);
  }
S
slzhou 已提交
1489 1490 1491
  if (newState.numOfResult == 1 || state.numOfResult == 1) {
    GET_RES_INFO(pCtx)->numOfRes = 1;
  }
S
shenglian zhou 已提交
1492

1493 1494
  blockDataDestroy(inputBlock);
  taosArrayDestroy(tempBlock.pDataBlock);
S
shenglian zhou 已提交
1495 1496

  taosMemoryFree(newState.buf);
1497
  return TSDB_CODE_SUCCESS;
S
shenglian zhou 已提交
1498 1499 1500
}

int32_t udfAggFinalize(struct SqlFunctionCtx *pCtx, SSDataBlock* pBlock) {
S
slzhou 已提交
1501
  SUdfAggRes* udfRes = (SUdfAggRes *)GET_ROWCELL_INTERBUF(GET_RES_INFO(pCtx));
S
slzhou 已提交
1502
  SClientUdfUvSession *session = udfRes->session;
S
slzhou 已提交
1503 1504
  udfRes->finalResBuf = (char*)udfRes + sizeof(SUdfAggRes);
  udfRes->interResBuf = (char*)udfRes + sizeof(SUdfAggRes) + session->outputLen;
S
shenglian zhou 已提交
1505 1506


S
slzhou 已提交
1507
  SUdfInterBuf resultBuf = {0};
S
slzhou 已提交
1508
  SUdfInterBuf state = {.buf = udfRes->interResBuf,
1509
                        .bufLen = session->bufSize,
S
slzhou 已提交
1510
                        .numOfResult = udfRes->interResNum};
1511 1512 1513 1514 1515 1516 1517 1518 1519 1520
  int32_t udfCallCode= 0;
  udfCallCode= callUdfAggFinalize(session, &state, &resultBuf);
  if (udfCallCode!= 0) {
    fnError("udfAggFinalize error. callUdfAggFinalize step. udf code:%d", udfCallCode);
    GET_RES_INFO(pCtx)->numOfRes = 0;
  } else {
    memcpy(udfRes->finalResBuf, resultBuf.buf, session->outputLen);
    udfRes->finalResNum = resultBuf.numOfResult;
    GET_RES_INFO(pCtx)->numOfRes = udfRes->finalResNum;
  }
S
shenglian zhou 已提交
1521

1522 1523 1524
  int32_t code = teardownUdf(session);
  if (code != 0) {
    fnError("udfAggFinalize error. teardownUdf step. udf code: %d", code);
S
slzhou 已提交
1525
  }
1526

1527
  return functionFinalizeWithResultBuf(pCtx, pBlock, udfRes->finalResBuf);
1528

S
shenglian zhou 已提交
1529
}