[Fix] Sockets not properly closed.
Sockets were not properly closed/released, causing the server to run out of connections or reaching the maxclient limit. Removed aeWinCloseSocket() in favor of the posix close() replacement. This commit fixes: - https://github.com/MSOpenTech/redis/issues/294 - https://github.com/MSOpenTech/redis/issues/282
This commit is contained in:
@@ -0,0 +1,32 @@
|
||||
/*
|
||||
* Copyright (c), Microsoft Open Technologies, Inc.
|
||||
* All rights reserved.
|
||||
* Redistribution and use in source and binary forms, with or without
|
||||
* modification, are permitted provided that the following conditions are met:
|
||||
* - Redistributions of source code must retain the above copyright notice,
|
||||
* this list of conditions and the following disclaimer.
|
||||
* - Redistributions in binary form must reproduce the above copyright notice,
|
||||
* this list of conditions and the following disclaimer in the documentation
|
||||
* and/or other materials provided with the distribution.
|
||||
* THIS SOFTWARE IS PROVIDED BY THE COPYRIGHT HOLDERS AND CONTRIBUTORS "AS IS"
|
||||
* AND ANY EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT LIMITED TO, THE
|
||||
* IMPLIED WARRANTIES OF MERCHANTABILITY AND FITNESS FOR A PARTICULAR PURPOSE ARE
|
||||
* DISCLAIMED. IN NO EVENT SHALL THE COPYRIGHT HOLDER OR CONTRIBUTORS BE LIABLE
|
||||
* FOR ANY DIRECT, INDIRECT, INCIDENTAL, SPECIAL, EXEMPLARY, OR CONSEQUENTIAL
|
||||
* DAMAGES (INCLUDING, BUT NOT LIMITED TO, PROCUREMENT OF SUBSTITUTE GOODS OR
|
||||
* SERVICES; LOSS OF USE, DATA, OR PROFITS; OR BUSINESS INTERRUPTION) HOWEVER
|
||||
* CAUSED AND ON ANY THEORY OF LIABILITY, WHETHER IN CONTRACT, STRICT LIABILITY,
|
||||
* OR TORT (INCLUDING NEGLIGENCE OR OTHERWISE) ARISING IN ANY WAY OUT OF THE USE
|
||||
* OF THIS SOFTWARE, EVEN IF ADVISED OF THE POSSIBILITY OF SUCH DAMAGE.
|
||||
*/
|
||||
|
||||
#ifndef WIN32_ASSERT_H
|
||||
#define WIN32_ASSERT_H
|
||||
|
||||
#ifdef _DEBUG
|
||||
#define ASSERT(condition) { if(!(condition)){ fprintf(stderr, "ASSERT FAILED: %s @ %s::%s (%d)\n", #condition , __FILE__, __FUNCTION__, __LINE__); DebugBreak();} }
|
||||
#else
|
||||
#define ASSERT(condition)
|
||||
#endif
|
||||
|
||||
#endif
|
||||
@@ -33,11 +33,16 @@
|
||||
#include "win32_util.h"
|
||||
#include "Win32_RedisLog.h"
|
||||
#include "Win32_Common.h"
|
||||
|
||||
#include "Win32_Assert.h"
|
||||
using namespace std;
|
||||
|
||||
#define CATCH_AND_REPORT() catch(const std::exception &){::redisLog(REDIS_WARNING, "FDAPI: std exception");}catch(...){::redisLog(REDIS_WARNING, "FDAPI: other exception");}
|
||||
|
||||
static fnIOCP_OnSocketClose* IOCP_OnSocketClose;
|
||||
void FDAPI_SetOnSocketClose(fnIOCP_OnSocketClose* iocpOnSocketClose) {
|
||||
IOCP_OnSocketClose = iocpOnSocketClose;
|
||||
}
|
||||
|
||||
extern "C" {
|
||||
// FD lookup Winsock equivalents for Win32_wsiocp.c
|
||||
redis_WSASetLastError WSASetLastError = NULL;
|
||||
@@ -308,30 +313,62 @@ int redis_pipe_impl(int *pfds) {
|
||||
return err;
|
||||
}
|
||||
|
||||
void FDAPI_ClearSocketInfo(int rfd) {
|
||||
SocketInfo* socketInfo = RFDMap::getInstance().lookupSocketInfo(rfd);
|
||||
|
||||
ASSERT(socketInfo != NULL)
|
||||
ASSERT(socketInfo->socket == INVALID_SOCKET)
|
||||
|
||||
if (socketInfo != NULL) {
|
||||
if (socketInfo->socket == INVALID_SOCKET) {
|
||||
RFDMap::getInstance().removeRFDToSocketInfo(rfd);
|
||||
return;
|
||||
} else {
|
||||
redisLog(REDIS_WARNING, "FDAPI_ClearSocketInfo called on non closed socket.");
|
||||
}
|
||||
} else {
|
||||
redisLog(REDIS_WARNING, "FDAPI_ClearSocketInfo called on non attached socket.");
|
||||
}
|
||||
}
|
||||
|
||||
// In unix a fd is a fd. All are closed with close().
|
||||
int redis_close_impl(RFD rfd) {
|
||||
try {
|
||||
SOCKET s = RFDMap::getInstance().lookupSocket(rfd);
|
||||
if( s != INVALID_SOCKET ) {
|
||||
RFDMap::getInstance().removeSocket(s);
|
||||
return f_closesocket(s);
|
||||
SocketInfo* socketInfo = RFDMap::getInstance().lookupSocketInfo(rfd);
|
||||
if (socketInfo != NULL) {
|
||||
|
||||
ASSERT(socketInfo->socket != INVALID_SOCKET)
|
||||
|
||||
if (socketInfo->socket != INVALID_SOCKET) {
|
||||
SOCKET socket = socketInfo->socket;
|
||||
socketInfo->socket = INVALID_SOCKET;
|
||||
|
||||
BOOL socketStateDeleted = FALSE;
|
||||
if (IOCP_OnSocketClose != NULL) {
|
||||
socketStateDeleted = IOCP_OnSocketClose(rfd);
|
||||
}
|
||||
|
||||
if (socketStateDeleted == TRUE) {
|
||||
RFDMap::getInstance().removeRFDToSocketInfo(rfd);
|
||||
}
|
||||
|
||||
RFDMap::getInstance().removeSocket(socket);
|
||||
return f_closesocket(socket);
|
||||
}
|
||||
} else {
|
||||
int posixFD = RFDMap::getInstance().lookupPosixFD(rfd);
|
||||
if(posixFD != -1) {
|
||||
if (posixFD != RFDMap::INVALID_FD) {
|
||||
RFDMap::getInstance().removePosixFD(posixFD);
|
||||
int retval = crt_close(posixFD);
|
||||
if( retval == -1 ) {
|
||||
if (retval == -1) {
|
||||
errno = GetLastError();
|
||||
}
|
||||
return retval;
|
||||
}
|
||||
else {
|
||||
errno = EBADF;
|
||||
return -1;
|
||||
}
|
||||
}
|
||||
} CATCH_AND_REPORT()
|
||||
|
||||
errno = EBADF;
|
||||
return -1;
|
||||
}
|
||||
|
||||
|
||||
@@ -190,6 +190,8 @@ typedef u_int64 (*redis_lseek64)(int fd, u_int64 offset, int whence);
|
||||
typedef intptr_t (*redis_get_osfhandle)(int fd);
|
||||
typedef int (*redis_open_osfhandle)(intptr_t osfhandle, int flags);
|
||||
|
||||
typedef BOOL fnIOCP_OnSocketClose(int rfd);
|
||||
|
||||
// access() mode definitions
|
||||
#define X_OK 0
|
||||
#define W_OK 2
|
||||
@@ -257,7 +259,9 @@ HANDLE FDAPI_CreateIoCompletionPortOnFD(int FD, HANDLE ExistingCompletionPort,
|
||||
BOOL FDAPI_AcceptEx(int listenFD,int acceptFD,PVOID lpOutputBuffer,DWORD dwReceiveDataLength,DWORD dwLocalAddressLength,DWORD dwRemoteAddressLength,LPDWORD lpdwBytesReceived,LPOVERLAPPED lpOverlapped);
|
||||
BOOL FDAPI_ConnectEx(int fd,const struct sockaddr *name,int namelen,PVOID lpSendBuffer,DWORD dwSendDataLength,LPDWORD lpdwBytesSent,LPOVERLAPPED lpOverlapped);
|
||||
void FDAPI_GetAcceptExSockaddrs(int fd, PVOID lpOutputBuffer,DWORD dwReceiveDataLength,DWORD dwLocalAddressLength,DWORD dwRemoteAddressLength,LPSOCKADDR *LocalSockaddr,LPINT LocalSockaddrLength,LPSOCKADDR *RemoteSockaddr,LPINT RemoteSockaddrLength);
|
||||
int FDAPI_UpdateAcceptContext( int fd );
|
||||
int FDAPI_UpdateAcceptContext(int fd);
|
||||
void FDAPI_SetOnSocketClose(fnIOCP_OnSocketClose *iocpOnSocketClose);
|
||||
void FDAPI_ClearSocketInfo(int fd);
|
||||
|
||||
// other networking functions
|
||||
BOOL ParseStorageAddress(const char *ip, int port, SOCKADDR_STORAGE* pSotrageAddr);
|
||||
|
||||
@@ -41,6 +41,7 @@
|
||||
<ItemGroup>
|
||||
<ClInclude Include="win32fixes.h" />
|
||||
<ClInclude Include="Win32_ANSI.h" />
|
||||
<ClInclude Include="Win32_Assert.h" />
|
||||
<ClInclude Include="Win32_CommandLine.h" />
|
||||
<ClInclude Include="Win32_dlmalloc.h" />
|
||||
<ClInclude Include="Win32_EventLog.h" />
|
||||
|
||||
@@ -20,7 +20,6 @@
|
||||
* OF THIS SOFTWARE, EVEN IF ADVISED OF THE POSSIBILITY OF SUCH DAMAGE.
|
||||
*/
|
||||
|
||||
|
||||
#include "win32_types.h"
|
||||
#include "win32_rfdmap.h"
|
||||
|
||||
@@ -46,7 +45,7 @@ RFD RFDMap::getNextRFDAvailable() {
|
||||
rfd = RFDRecyclePool.front();
|
||||
RFDRecyclePool.pop();
|
||||
} else {
|
||||
rfd = (int) SocketToRFDMap.size() + (int) PosixFDToRFDMap.size();
|
||||
rfd = RFDMap::next_available_rfd++;
|
||||
}
|
||||
LeaveCriticalSection(&mutex);
|
||||
return rfd;
|
||||
@@ -60,7 +59,10 @@ RFD RFDMap::addSocket(SOCKET s) {
|
||||
} else {
|
||||
rfd = getNextRFDAvailable();
|
||||
SocketToRFDMap[s] = rfd;
|
||||
RFDToSocketMap[rfd] = s;
|
||||
|
||||
SocketInfo socket_info;
|
||||
socket_info.socket = s;
|
||||
RFDToSocketInfoMap[rfd] = socket_info;
|
||||
}
|
||||
LeaveCriticalSection(&mutex);
|
||||
return rfd;
|
||||
@@ -68,13 +70,14 @@ RFD RFDMap::addSocket(SOCKET s) {
|
||||
|
||||
void RFDMap::removeSocket(SOCKET s) {
|
||||
EnterCriticalSection(&mutex);
|
||||
map<SOCKET, RFD>::iterator mit = SocketToRFDMap.find(s);
|
||||
if (mit != SocketToRFDMap.end()) {
|
||||
RFD rfd = (*mit).second;
|
||||
RFDRecyclePool.push(rfd);
|
||||
RFDToSocketMap.erase(rfd);
|
||||
SocketToRFDMap.erase(s);
|
||||
}
|
||||
SocketToRFDMap.erase(s);
|
||||
LeaveCriticalSection(&mutex);
|
||||
}
|
||||
|
||||
void RFDMap::removeRFDToSocketInfo(RFD rfd) {
|
||||
EnterCriticalSection(&mutex);
|
||||
RFDToSocketInfoMap.erase(rfd);
|
||||
RFDRecyclePool.push(rfd);
|
||||
LeaveCriticalSection(&mutex);
|
||||
}
|
||||
|
||||
@@ -93,9 +96,9 @@ RFD RFDMap::addPosixFD(int posixFD) {
|
||||
}
|
||||
|
||||
void RFDMap::removePosixFD(int posixFD) {
|
||||
// posixFD between 0 and 2 should never be removed since they are assigned
|
||||
// to stdin, stdout and stderr
|
||||
if (posixFD > 2) {
|
||||
// posixFD between RESERVED_RFD_INDEX_START and RESERVED_RFD_INDEX_END
|
||||
// should never be removed.
|
||||
if (posixFD > RFDMap::RESERVED_RFD_INDEX_END) {
|
||||
EnterCriticalSection(&mutex);
|
||||
map<int, RFD>::iterator mit = PosixFDToRFDMap.find(posixFD);
|
||||
if (mit != PosixFDToRFDMap.end()) {
|
||||
@@ -111,19 +114,30 @@ void RFDMap::removePosixFD(int posixFD) {
|
||||
SOCKET RFDMap::lookupSocket(RFD rfd) {
|
||||
SOCKET socket = INVALID_SOCKET;
|
||||
EnterCriticalSection(&mutex);
|
||||
if (RFDToSocketMap.find(rfd) != RFDToSocketMap.end()) {
|
||||
socket = RFDToSocketMap[rfd];
|
||||
if (RFDToSocketInfoMap.find(rfd) != RFDToSocketInfoMap.end()) {
|
||||
socket = RFDToSocketInfoMap[rfd].socket;
|
||||
}
|
||||
LeaveCriticalSection(&mutex);
|
||||
return socket;
|
||||
}
|
||||
|
||||
SocketInfo* RFDMap::lookupSocketInfo(RFD rfd) {
|
||||
SocketInfo* socket_info = NULL;
|
||||
EnterCriticalSection(&mutex);
|
||||
if (RFDToSocketInfoMap.find(rfd) != RFDToSocketInfoMap.end()) {
|
||||
socket_info = &RFDToSocketInfoMap[rfd];
|
||||
}
|
||||
LeaveCriticalSection(&mutex);
|
||||
return socket_info;
|
||||
}
|
||||
|
||||
int RFDMap::lookupPosixFD(RFD rfd) {
|
||||
int posixFD = -1;
|
||||
int posixFD = RFDMap::INVALID_FD;
|
||||
EnterCriticalSection(&mutex);
|
||||
if (RFDToPosixFDMap.find(rfd) != RFDToPosixFDMap.end()) {
|
||||
posixFD = RFDToPosixFDMap[rfd];
|
||||
} else if (rfd >= 0 && rfd <= 2) {
|
||||
} else if (rfd >= RFDMap::RESERVED_RFD_INDEX_START
|
||||
&& rfd <= RFDMap::RESERVED_RFD_INDEX_END) {
|
||||
posixFD = rfd;
|
||||
}
|
||||
LeaveCriticalSection(&mutex);
|
||||
|
||||
@@ -33,6 +33,10 @@ using namespace std;
|
||||
|
||||
typedef int RFD; // Redis File Descriptor
|
||||
|
||||
typedef struct {
|
||||
SOCKET socket;
|
||||
} SocketInfo;
|
||||
|
||||
/* In UNIX File Descriptors increment by one for each new one. Windows handles
|
||||
* do not follow the same rule. Additionally UNIX uses a 32-bit int to
|
||||
* represent a FD while Windows_x64 uses a 64-bit value to represent a handle.
|
||||
@@ -41,7 +45,7 @@ typedef int RFD; // Redis File Descriptor
|
||||
* should be treated as an opaque value and not be cast to a 32-bit int. In
|
||||
* order to not break existing code that relies on the maximum FD value to
|
||||
* indicate the number of handles that have been created (and other UNIXisms),
|
||||
* this code maps SOCKET handles to a virtual FD number starting at 3 (0,1 and
|
||||
* this code maps SOCKET handles to a virtual FD number starting at 3 (0, 1 and
|
||||
* 2 are reserved for stdin, stdout and stderr).
|
||||
*/
|
||||
class RFDMap {
|
||||
@@ -54,50 +58,88 @@ private:
|
||||
void operator=(RFDMap const&); // Don't implement to guarantee singleton semantics
|
||||
|
||||
private:
|
||||
map<SOCKET, RFD> SocketToRFDMap;
|
||||
map<SOCKET, int> SocketToFlagsMap;
|
||||
map<int, RFD> PosixFDToRFDMap;
|
||||
map<RFD, int> RFDToPosixFDMap;
|
||||
map<RFD, SOCKET> RFDToSocketMap;
|
||||
queue<RFD> RFDRecyclePool;
|
||||
map<SOCKET, RFD> SocketToRFDMap;
|
||||
map<SOCKET, int> SocketToFlagsMap;
|
||||
map<int, RFD> PosixFDToRFDMap;
|
||||
map<RFD, SocketInfo> RFDToSocketInfoMap;
|
||||
map<RFD, int> RFDToPosixFDMap;
|
||||
queue<RFD> RFDRecyclePool;
|
||||
|
||||
private:
|
||||
CRITICAL_SECTION mutex;
|
||||
|
||||
public:
|
||||
const static int INVALID_RFD = -1;
|
||||
const static int INVALID_FD = -1;
|
||||
|
||||
private:
|
||||
/* Gets the next available Redis File Descriptor. Redis File Descriptors are always
|
||||
non-negative integers, with the first three being reserved for stdin(0),
|
||||
stdout(1) and stderr(2). */
|
||||
/*
|
||||
Gets the next available Redis File Descriptor. RFDs are always
|
||||
non-negative integers, with the first three being reserved for
|
||||
stdin(0), stdout(1) and stderr(2).
|
||||
*/
|
||||
RFD getNextRFDAvailable();
|
||||
|
||||
const static int RESERVED_RFD_INDEX_START = 0;
|
||||
const static int RESERVED_RFD_INDEX_END = 2;
|
||||
int next_available_rfd = RESERVED_RFD_INDEX_END + 1;
|
||||
|
||||
public:
|
||||
/* Adds a socket to the socket map. Returns the redis file descriptor value for
|
||||
the socket. Returns invalidRFD if the socket is already added to the
|
||||
collection. */
|
||||
/*
|
||||
Adds a socket to SocketToRFDMap and to RFDToSocketInfoMap.
|
||||
Returns the RFD value for the socket.
|
||||
Returns INVALID_RFD if the socket is already added to the collection.
|
||||
*/
|
||||
RFD addSocket(SOCKET s);
|
||||
|
||||
/* Removes a socket from the list of sockets. Also removes the associated
|
||||
file descriptor. */
|
||||
/* Removes a socket from SocketToRFDMap. */
|
||||
void removeSocket(SOCKET s);
|
||||
|
||||
/* Adds a posixFD (used with low-level CRT posix file functions) to the posixFD map. Returns
|
||||
the redis file descriptor value for the posixFD. Returns invalidRFD if the posicFD is already
|
||||
added to the collection. */
|
||||
/*
|
||||
Removes a RFD from RFDToSocketInfoMap.
|
||||
It frees the associated RFD adding it to RFDRecyclePool.
|
||||
*/
|
||||
void removeRFDToSocketInfo(RFD rfd);
|
||||
|
||||
/*
|
||||
Adds a posixFD (used with low-level CRT posix file functions) to RFDToPosixFDMap.
|
||||
Returns the RFD value for the posixFD.
|
||||
Returns the existing RFD if the posixFD is already present in the collection.
|
||||
*/
|
||||
RFD addPosixFD(int posixFD);
|
||||
|
||||
/* Removes a socket from the list of sockets. Also removes the associated
|
||||
file descriptor. */
|
||||
/*
|
||||
Removes a socket from RFDToPosixFDMap.
|
||||
It frees the associated RFD adding it to RFDRecyclePool.
|
||||
*/
|
||||
void removePosixFD(int posixFD);
|
||||
|
||||
/* Returns the socket associated with a file descriptor. */
|
||||
/*
|
||||
Returns the socket associated with a RFD.
|
||||
Returns INVALID_SOCKET if the socket is not found.
|
||||
*/
|
||||
SOCKET lookupSocket(RFD rfd);
|
||||
|
||||
/* Returns the socket associated with a file descriptor. */
|
||||
/*
|
||||
Returns a pointer to the socket info structure associated with a RFD.
|
||||
Return NULL if the info socket info structure is not found.
|
||||
*/
|
||||
SocketInfo* lookupSocketInfo(RFD rfd);
|
||||
|
||||
/*
|
||||
Returns the posixFD associated with a RFD.
|
||||
Returns INVALID_FD if the posixFD is not found.
|
||||
*/
|
||||
int lookupPosixFD(RFD rfd);
|
||||
|
||||
/*
|
||||
Sets the socket flags.
|
||||
*/
|
||||
void SetSocketFlags(SOCKET s, int flags);
|
||||
|
||||
/*
|
||||
Returns the socket flags.
|
||||
Returns 0 if the flags are not found.
|
||||
*/
|
||||
int GetSocketFlags(SOCKET s);
|
||||
};
|
||||
|
||||
@@ -31,6 +31,7 @@
|
||||
static void *iocpState;
|
||||
static HANDLE iocph;
|
||||
static fnGetSockState * aeGetSockState;
|
||||
static fnGetSockState * aeGetExistingSockState;
|
||||
static fnDelSockState * aeDelSockState;
|
||||
|
||||
#define SUCCEEDED_WITH_IOCP(result) ((result) || (GetLastError() == ERROR_IO_PENDING))
|
||||
@@ -466,37 +467,31 @@ void aeShutdown(int fd) {
|
||||
}
|
||||
}
|
||||
|
||||
/* when closing socket, need to unassociate completion port */
|
||||
int aeWinCloseSocket(int fd) {
|
||||
aeSockState *sockstate;
|
||||
|
||||
if ((sockstate = aeGetSockState(iocpState, fd)) == NULL) {
|
||||
close(fd);
|
||||
return 0;
|
||||
BOOL WSIOCP_OnSocketClose(int rfd) {
|
||||
aeSockState *socketState;
|
||||
if ((socketState = aeGetExistingSockState(iocpState, rfd)) == NULL) {
|
||||
return TRUE;
|
||||
}
|
||||
|
||||
aeShutdown(fd);
|
||||
sockstate->masks &= ~(SOCKET_ATTACHED | AE_WRITABLE | AE_READABLE);
|
||||
|
||||
if (sockstate->wreqs == 0 &&
|
||||
(sockstate->masks & (READ_QUEUED | CONNECT_PENDING | SOCKET_ATTACHED)) == 0) {
|
||||
close(fd);
|
||||
} else {
|
||||
sockstate->masks |= CLOSE_PENDING;
|
||||
socketState->masks &= ~(SOCKET_ATTACHED | AE_WRITABLE | AE_READABLE);
|
||||
if (socketState->wreqs != 0 || (socketState->masks & (READ_QUEUED | CONNECT_PENDING)) != 0) {
|
||||
socketState->masks |= CLOSE_PENDING;
|
||||
}
|
||||
aeDelSockState(iocpState, sockstate);
|
||||
|
||||
return 0;
|
||||
return aeDelSockState(iocpState, socketState);
|
||||
}
|
||||
|
||||
void aeWinInit(void *state,
|
||||
HANDLE iocp,
|
||||
fnGetSockState *getSockState,
|
||||
fnGetSockState *getExistingSockState,
|
||||
fnDelSockState *delSockState) {
|
||||
iocpState = state;
|
||||
iocph = iocp;
|
||||
aeGetSockState = getSockState;
|
||||
aeGetExistingSockState = getExistingSockState;
|
||||
aeDelSockState = delSockState;
|
||||
FDAPI_SetOnSocketClose(WSIOCP_OnSocketClose);
|
||||
}
|
||||
|
||||
void aeWinCleanup() {
|
||||
|
||||
@@ -59,7 +59,7 @@ typedef struct aeSockState {
|
||||
} aeSockState;
|
||||
|
||||
typedef aeSockState * fnGetSockState(void *apistate, int fd);
|
||||
typedef void fnDelSockState(void *apistate, aeSockState *sockState);
|
||||
typedef BOOL fnDelSockState(void *apistate, aeSockState *sockState);
|
||||
|
||||
#define READ_QUEUED 0x000100
|
||||
#define SOCKET_ATTACHED 0x000400
|
||||
@@ -68,7 +68,7 @@ typedef void fnDelSockState(void *apistate, aeSockState *sockState);
|
||||
#define CONNECT_PENDING 0x002000
|
||||
#define CLOSE_PENDING 0x004000
|
||||
|
||||
void aeWinInit(void *state, HANDLE iocp, fnGetSockState *getSockState, fnDelSockState *delSockState);
|
||||
void aeWinInit(void *state, HANDLE iocp, fnGetSockState *getSockState, fnGetSockState *getExistingSockState, fnDelSockState *delSockState);
|
||||
void aeWinCleanup();
|
||||
|
||||
void* CallocMemoryNoCOW(size_t size);
|
||||
|
||||
@@ -292,7 +292,6 @@ typedef struct aeWinSendReq {
|
||||
|
||||
|
||||
int aeWinSocketAttach(int fd);
|
||||
int aeWinCloseSocket(int fd);
|
||||
int aeWinReceiveDone(int fd);
|
||||
int aeWinSocketSend(int fd, char *buf, int len,
|
||||
void *eventLoop, void *client, void *data, void *proc);
|
||||
|
||||
+37
-46
@@ -61,42 +61,6 @@ int aeSocketIndex(int fd) {
|
||||
return fd;
|
||||
}
|
||||
|
||||
/* get data for socket / fd being monitored. Create if not found*/
|
||||
aeSockState *aeGetSockState(void *apistate, int fd) {
|
||||
int sindex;
|
||||
listNode *node;
|
||||
list *socklist;
|
||||
aeSockState *sockState;
|
||||
if (apistate == NULL) return NULL;
|
||||
|
||||
sindex = aeSocketIndex(fd);
|
||||
socklist = &(((aeApiState *)apistate)->lookup[sindex]);
|
||||
node = listFirst(socklist);
|
||||
while (node != NULL) {
|
||||
sockState = (aeSockState *)listNodeValue(node);
|
||||
if (sockState->fd == fd) {
|
||||
return sockState;
|
||||
}
|
||||
node = listNextNode(node);
|
||||
}
|
||||
// not found. Do lazy create of sockState.
|
||||
sockState = (aeSockState *) CallocMemoryNoCOW(sizeof(aeSockState));
|
||||
if (sockState != NULL) {
|
||||
sockState->fd = fd;
|
||||
sockState->masks = 0;
|
||||
sockState->wreqs = 0;
|
||||
sockState->reqs = NULL;
|
||||
memset(&sockState->wreqlist, 0, sizeof(sockState->wreqlist));
|
||||
|
||||
if (listAddNodeHead(socklist, sockState) != NULL) {
|
||||
return sockState;
|
||||
} else {
|
||||
FreeMemoryNoCOW(sockState);
|
||||
}
|
||||
}
|
||||
return NULL;
|
||||
}
|
||||
|
||||
/* get data for socket / fd being monitored */
|
||||
aeSockState *aeGetExistingSockState(void *apistate, int fd) {
|
||||
int sindex;
|
||||
@@ -119,6 +83,30 @@ aeSockState *aeGetExistingSockState(void *apistate, int fd) {
|
||||
return NULL;
|
||||
}
|
||||
|
||||
/* get data for socket / fd being monitored. Create if not found*/
|
||||
aeSockState *aeGetSockState(void *apistate, int fd) {
|
||||
aeSockState *sockState = aeGetExistingSockState(apistate, fd);
|
||||
if (sockState == NULL) {
|
||||
// Not found. Do lazy create of sockState.
|
||||
sockState = (aeSockState *) CallocMemoryNoCOW(sizeof(aeSockState));
|
||||
if (sockState != NULL) {
|
||||
sockState->fd = fd;
|
||||
sockState->masks = 0;
|
||||
sockState->wreqs = 0;
|
||||
sockState->reqs = NULL;
|
||||
memset(&sockState->wreqlist, 0, sizeof(sockState->wreqlist));
|
||||
|
||||
int sindex = aeSocketIndex(fd);
|
||||
list *socklist = &(((aeApiState *) apistate)->lookup[sindex]);
|
||||
if (listAddNodeHead(socklist, sockState) == NULL) {
|
||||
FreeMemoryNoCOW(sockState);
|
||||
return NULL;
|
||||
}
|
||||
}
|
||||
}
|
||||
return sockState;
|
||||
}
|
||||
|
||||
// find matching value in list and remove. If found return 1
|
||||
int removeMatchFromList(list *socklist, void *value) {
|
||||
listNode *node;
|
||||
@@ -136,12 +124,12 @@ int removeMatchFromList(list *socklist, void *value) {
|
||||
|
||||
/* delete data for socket / fd being monitored
|
||||
or move to the closing queue if operations are pending.
|
||||
Return 1 if deleted or not found, 0 if pending*/
|
||||
void aeDelSockState(void *apistate, aeSockState *sockState) {
|
||||
Return TRUE if deleted or not found, FLASE if pending*/
|
||||
BOOL aeDelSockState(void *apistate, aeSockState *sockState) {
|
||||
int sindex;
|
||||
list *socklist;
|
||||
|
||||
if (apistate == NULL) return;
|
||||
if (apistate == NULL) return TRUE;
|
||||
|
||||
if (sockState->wreqs == 0 &&
|
||||
(sockState->masks & (READ_QUEUED | CONNECT_PENDING | SOCKET_ATTACHED | CLOSE_PENDING)) == 0) {
|
||||
@@ -150,13 +138,13 @@ void aeDelSockState(void *apistate, aeSockState *sockState) {
|
||||
socklist = &(((aeApiState *)apistate)->lookup[sindex]);
|
||||
if (removeMatchFromList(socklist, sockState) == 1) {
|
||||
FreeMemoryNoCOW(sockState);
|
||||
return;
|
||||
return TRUE;
|
||||
}
|
||||
// try closing list
|
||||
socklist = &(((aeApiState *)apistate)->closing);
|
||||
if (removeMatchFromList(socklist, sockState) == 1) {
|
||||
FreeMemoryNoCOW(sockState);
|
||||
return;
|
||||
return TRUE;
|
||||
}
|
||||
} else {
|
||||
// not safe to delete. Move to closing
|
||||
@@ -168,6 +156,8 @@ void aeDelSockState(void *apistate, aeSockState *sockState) {
|
||||
listAddNodeHead(socklist, sockState);
|
||||
}
|
||||
}
|
||||
|
||||
return FALSE;
|
||||
}
|
||||
|
||||
/* Called by ae to initialize state */
|
||||
@@ -198,7 +188,7 @@ static int aeApiCreate(aeEventLoop *eventLoop) {
|
||||
state->setsize = eventLoop->setsize;
|
||||
eventLoop->apidata = state;
|
||||
/* initialize the IOCP socket code with state reference */
|
||||
aeWinInit(state, state->iocp, aeGetSockState, aeDelSockState);
|
||||
aeWinInit(state, state->iocp, aeGetSockState, aeGetExistingSockState, aeDelSockState);
|
||||
return 0;
|
||||
}
|
||||
|
||||
@@ -414,12 +404,13 @@ static int aeApiPoll(aeEventLoop *eventLoop, struct timeval *tvp) {
|
||||
}
|
||||
if (sockstate->wreqs == 0 &&
|
||||
(sockstate->masks & (CONNECT_PENDING | READ_QUEUED | SOCKET_ATTACHED)) == 0) {
|
||||
if ((sockstate->masks & CLOSE_PENDING) != 0) {
|
||||
close(rfd);
|
||||
if (sockstate->masks & CLOSE_PENDING) {
|
||||
sockstate->masks &= ~(CLOSE_PENDING);
|
||||
}
|
||||
// safe to delete sockstate
|
||||
aeDelSockState(state, sockstate);
|
||||
if (aeDelSockState(state, sockstate)) {
|
||||
FDAPI_ClearSocketInfo(rfd);
|
||||
}
|
||||
sockstate = NULL;
|
||||
}
|
||||
break;
|
||||
}
|
||||
|
||||
+2
-18
@@ -159,11 +159,7 @@ void migrateCommand(redisClient *c) {
|
||||
return;
|
||||
}
|
||||
if ((aeWait(fd,AE_WRITABLE,timeout) & AE_WRITABLE) == 0) {
|
||||
#ifdef WIN32_IOCP
|
||||
aeWinCloseSocket(fd);
|
||||
#else
|
||||
close(fd);
|
||||
#endif
|
||||
addReplySds(c,sdsnew("-IOERR error or timeout connecting to the client\r\n"));
|
||||
return;
|
||||
}
|
||||
@@ -257,21 +253,13 @@ void migrateCommand(redisClient *c) {
|
||||
}
|
||||
|
||||
sdsfree(cmd.io.buffer.ptr);
|
||||
#ifdef WIN32_IOCP
|
||||
aeWinCloseSocket(fd);
|
||||
#else
|
||||
close(fd);
|
||||
#endif
|
||||
close(fd);
|
||||
return;
|
||||
|
||||
socket_wr_err:
|
||||
addReplySds(c,sdsnew("-IOERR error or timeout writing to target instance\r\n"));
|
||||
sdsfree(cmd.io.buffer.ptr);
|
||||
#ifdef WIN32_IOCP
|
||||
aeWinCloseSocket(fd);
|
||||
#else
|
||||
close(fd);
|
||||
#endif
|
||||
return;
|
||||
|
||||
socket_rd_err:
|
||||
@@ -280,10 +268,6 @@ socket_rd_err:
|
||||
#endif
|
||||
addReplySds(c,sdsnew("-IOERR error or timeout reading from target node\r\n"));
|
||||
sdsfree(cmd.io.buffer.ptr);
|
||||
#ifdef WIN32_IOCP
|
||||
aeWinCloseSocket(fd);
|
||||
#else
|
||||
close(fd);
|
||||
#endif
|
||||
close(fd);
|
||||
return;
|
||||
}
|
||||
|
||||
+1
-14
@@ -72,11 +72,7 @@ redisClient *createClient(int fd) {
|
||||
if (aeCreateFileEvent(server.el,fd,AE_READABLE,
|
||||
readQueryFromClient, c) == AE_ERR)
|
||||
{
|
||||
#ifdef WIN32_IOCP
|
||||
aeWinCloseSocket(fd);
|
||||
#else
|
||||
close(fd);
|
||||
#endif
|
||||
close(fd);
|
||||
zfree(c);
|
||||
return NULL;
|
||||
}
|
||||
@@ -557,11 +553,7 @@ static void acceptCommonHandler(int fd, int flags) {
|
||||
redisLog(REDIS_WARNING,
|
||||
"Error registering fd event for the new client: %s (fd=%d)",
|
||||
strerror(errno),fd);
|
||||
#ifdef WIN32_IOCP
|
||||
aeWinCloseSocket(fd); /* May be already closed, just ingore errors */
|
||||
#else
|
||||
close(fd); /* May be already closed, just ingore errors */
|
||||
#endif
|
||||
return;
|
||||
}
|
||||
/* If maxclient directive is set and this is one client more... close the
|
||||
@@ -718,11 +710,6 @@ void freeClient(redisClient *c) {
|
||||
}
|
||||
listRelease(c->reply);
|
||||
freeClientArgv(c);
|
||||
#ifdef WIN32_IOCP
|
||||
aeWinCloseSocket(c->fd);
|
||||
#else
|
||||
|
||||
#endif
|
||||
/* Remove from the list of clients */
|
||||
if (c->fd != -1) {
|
||||
ln = listSearchKey(server.clients,c);
|
||||
|
||||
@@ -142,10 +142,6 @@ static void freeClient(client c) {
|
||||
listNode *ln;
|
||||
aeDeleteFileEvent(config.el,(int)c->context->fd,AE_WRITABLE);
|
||||
aeDeleteFileEvent(config.el,(int)c->context->fd,AE_READABLE);
|
||||
#ifdef WIN32_IOCP
|
||||
aeWinCloseSocket((int)c->context->fd);
|
||||
c->context->fd = 0;
|
||||
#endif
|
||||
redisFree(c->context);
|
||||
sdsfree(c->obuf);
|
||||
zfree(c->randptr);
|
||||
|
||||
+2
-19
@@ -924,20 +924,15 @@ void replicationAbortSyncTransfer(void) {
|
||||
redisAssert(server.repl_state == REDIS_REPL_TRANSFER);
|
||||
|
||||
aeDeleteFileEvent(server.el,server.repl_transfer_s,AE_READABLE);
|
||||
close(server.repl_transfer_s);
|
||||
close(server.repl_transfer_fd);
|
||||
#ifdef WIN32_IOCP
|
||||
aeWinCloseSocket(server.repl_transfer_s);
|
||||
if (server.repl_transfer_fd != -1) {
|
||||
close(server.repl_transfer_fd);
|
||||
server.repl_transfer_fd = -1;
|
||||
}
|
||||
if (server.repl_transfer_tmpfile != NULL) {
|
||||
unlink(server.repl_transfer_tmpfile);
|
||||
zfree(server.repl_transfer_tmpfile);
|
||||
server.repl_transfer_tmpfile = NULL;
|
||||
}
|
||||
#else
|
||||
close(server.repl_transfer_s);
|
||||
close(server.repl_transfer_fd);
|
||||
unlink(server.repl_transfer_tmpfile);
|
||||
zfree(server.repl_transfer_tmpfile);
|
||||
#endif
|
||||
@@ -1522,11 +1517,7 @@ void syncWithMaster(aeEventLoop *el, int fd, void *privdata, int mask) {
|
||||
return;
|
||||
|
||||
error:
|
||||
#ifdef WIN32_IOCP
|
||||
aeWinCloseSocket(fd);
|
||||
#else
|
||||
close(fd);
|
||||
#endif
|
||||
server.repl_transfer_s = -1;
|
||||
server.repl_state = REDIS_REPL_CONNECT;
|
||||
return;
|
||||
@@ -1545,11 +1536,7 @@ int connectWithMaster(void) {
|
||||
if (aeCreateFileEvent(server.el,fd,AE_READABLE|AE_WRITABLE,syncWithMaster,NULL) ==
|
||||
AE_ERR)
|
||||
{
|
||||
#ifdef WIN32_IOCP
|
||||
aeWinCloseSocket(fd);
|
||||
#else
|
||||
close(fd);
|
||||
#endif
|
||||
redisLog(REDIS_WARNING,"Can't create readable event for SYNC");
|
||||
return REDIS_ERR;
|
||||
}
|
||||
@@ -1568,11 +1555,7 @@ void undoConnectWithMaster(void) {
|
||||
redisAssert(server.repl_state == REDIS_REPL_CONNECTING ||
|
||||
server.repl_state == REDIS_REPL_RECEIVE_PONG);
|
||||
aeDeleteFileEvent(server.el,fd,AE_READABLE|AE_WRITABLE);
|
||||
#ifdef WIN32_IOCP
|
||||
aeWinCloseSocket(fd);
|
||||
#else
|
||||
close(fd);
|
||||
#endif
|
||||
server.repl_transfer_s = -1;
|
||||
server.repl_state = REDIS_REPL_CONNECT;
|
||||
}
|
||||
|
||||
@@ -331,10 +331,6 @@ static void redisAeCleanup(void *privdata) {
|
||||
redisAeEvents *e = (redisAeEvents*)privdata;
|
||||
redisAeDelRead(privdata);
|
||||
redisAeDelWrite(privdata);
|
||||
#ifdef WIN32_IOCP
|
||||
aeWinCloseSocket((int)e->fd);
|
||||
e->fd = 0;
|
||||
#endif
|
||||
zfree(e);
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user