--- libaitrpc/src/srv.c 2013/11/14 23:38:41 1.21.2.2 +++ libaitrpc/src/srv.c 2014/01/28 13:56:25 1.22.6.1 @@ -3,7 +3,7 @@ * by Michael Pounov * * $Author: misho $ -* $Id: srv.c,v 1.21.2.2 2013/11/14 23:38:41 misho Exp $ +* $Id: srv.c,v 1.22.6.1 2014/01/28 13:56:25 misho Exp $ * ************************************************************************** The ELWIX and AITNET software is distributed under the following @@ -66,7 +66,11 @@ static sched_task_func_t cbProto[SOCK_RAW + 1][4] = { { NULL, NULL, NULL, NULL } /* SOCK_RAW */ }; +/* Global Signal Argument when kqueue support disabled */ +static volatile uintptr_t _glSigArg = 0; + + void rpc_freeCli(rpc_cli_t * __restrict c) { @@ -223,7 +227,7 @@ txPacket(sched_task_t *task) if (ret) LOGERR; else - rpc_SetErr(ETIMEDOUT, "Timeout reached! Server not respond"); + rpc_SetErr(ETIMEDOUT, "Timeout reached! Client not respond"); /* close connection */ schedEvent(TASK_ROOT(task), cbProto[s->srv_proto][CB_CLOSECLIENT], TASK_ARG(task), 0, NULL, 0); @@ -346,7 +350,7 @@ rxPacket(sched_task_t *task) if (rlen) LOGERR; else - rpc_SetErr(ETIMEDOUT, "Timeout reached! Server not respond"); + rpc_SetErr(ETIMEDOUT, "Timeout reached! Client not respond"); schedEvent(TASK_ROOT(task), cbProto[s->srv_proto][CB_CLOSECLIENT], TASK_ARG(task), 0, NULL, 0); return NULL; @@ -448,7 +452,7 @@ txUDPPacket(sched_task_t *task) rpc_func_t *f = NULL; u_char *buf = AIT_GET_BUF(&c->cli_buf); struct tagRPCCall *rpc = (struct tagRPCCall*) buf; - int ret, wlen = sizeof(struct tagRPCCall); + int ret, estlen, wlen = sizeof(struct tagRPCCall); struct timespec ts = { DEF_RPC_TIMEOUT, 0 }; struct pollfd pfd; @@ -464,6 +468,13 @@ txUDPPacket(sched_task_t *task) rpc->call_rep.ret = RPC_ERROR(-1); rpc->call_rep.eno = RPC_ERROR(rpc_Errno); } else { + /* calc estimated length */ + estlen = ait_resideVars(RPC_RETVARS(c)) + wlen; + if (estlen > AIT_LEN(&c->cli_buf)) + AIT_RE_BUF(&c->cli_buf, estlen); + buf = AIT_GET_BUF(&c->cli_buf); + rpc = (struct tagRPCCall*) buf; + rpc->call_argc = htons(array_Size(RPC_RETVARS(c))); /* Go Encapsulate variables */ ret = ait_vars2buffer(buf + wlen, AIT_LEN(&c->cli_buf) - wlen, @@ -495,7 +506,7 @@ txUDPPacket(sched_task_t *task) if (ret) LOGERR; else - rpc_SetErr(ETIMEDOUT, "Timeout reached! Server not respond"); + rpc_SetErr(ETIMEDOUT, "Timeout reached! Client not respond"); /* close connection */ schedEvent(TASK_ROOT(task), cbProto[s->srv_proto][CB_CLOSECLIENT], TASK_ARG(task), 0, NULL, 0); @@ -537,9 +548,11 @@ rxUDPPacket(sched_task_t *task) } c = _allocClient(srv, &sa); - if (!c) + if (!c) { + EVERBOSE(1, "RPC client quota exceeded! Connection will be shutdown!\n"); + usleep(2000); /* blocked client delay */ goto end; - else { + } else { estlen = ntohl(rpc->call_len); if (estlen > AIT_LEN(&c->cli_buf)) AIT_RE_BUF(&c->cli_buf, estlen); @@ -566,11 +579,12 @@ rxUDPPacket(sched_task_t *task) if (rlen) LOGERR; else - rpc_SetErr(ETIMEDOUT, "Timeout reached! Server not respond"); + rpc_SetErr(ETIMEDOUT, "Timeout reached! Client not respond"); schedEvent(TASK_ROOT(task), cbProto[srv->srv_proto][CB_CLOSECLIENT], c, 0, NULL, 0); return NULL; } + salen = sa.ss.ss_len = sizeof(sockaddr_t); rlen = recvfrom(TASK_FD(task), buf, len, 0, &sa.sa, &salen); if (rlen == -1) { /* close connection */ @@ -578,6 +592,8 @@ rxUDPPacket(sched_task_t *task) c, 0, NULL, 0); return NULL; } + if (e_addrcmp(&c->cli_sa, &sa, 42)) + rlen ^= rlen; /* skip if arrive from different address */ } len = estlen; @@ -748,7 +764,8 @@ end: static void * flushBLOB(sched_task_t *task) { - rpc_srv_t *srv = TASK_ARG(task); + uintptr_t sigArg = atomic_load_acq_ptr(&_glSigArg); + rpc_srv_t *srv = sigArg ? (void*) sigArg : TASK_ARG(task); rpc_blob_t *b, *tmp; TAILQ_FOREACH_SAFE(b, &srv->srv_blob.blobs, blob_node, tmp) { @@ -758,7 +775,17 @@ flushBLOB(sched_task_t *task) e_free(b); } - schedSignalSelf(task); + if (!schedSignalSelf(task)) { + /* disabled kqueue support in libaitsched */ + struct sigaction sa; + + memset(&sa, 0, sizeof sa); + sigemptyset(&sa.sa_mask); + sa.sa_handler = (void (*)(int)) flushBLOB; + sa.sa_flags = SA_RESTART | SA_RESETHAND; + sigaction(SIGFBLOB, &sa, NULL); + } + return NULL; } @@ -972,7 +999,19 @@ rpc_srv_loopBLOBServer(rpc_srv_t * __restrict srv) return -1; } - schedSignal(srv->srv_blob.root, flushBLOB, srv, SIGFBLOB, NULL, 0); + if (!schedSignal(srv->srv_blob.root, flushBLOB, srv, SIGFBLOB, NULL, 0)) { + /* disabled kqueue support in libaitsched */ + struct sigaction sa; + + atomic_store_rel_ptr(&_glSigArg, (uintptr_t) srv); + + memset(&sa, 0, sizeof sa); + sigemptyset(&sa.sa_mask); + sa.sa_handler = (void (*)(int)) flushBLOB; + sa.sa_flags = SA_RESTART | SA_RESETHAND; + sigaction(SIGFBLOB, &sa, NULL); + } + if (!schedRead(srv->srv_blob.root, acceptBLOBClients, srv, srv->srv_blob.server.cli_sock, NULL, 0)) { rpc_SetErr(sched_GetErrno(), "%s", sched_GetError());