Networking: fixes some issues and memory leaks

This commit is contained in:
李通洲
2026-01-06 15:33:12 +08:00
parent 6ded0f9bb0
commit 1be90084ab
4 changed files with 103 additions and 64 deletions
+10 -10
View File
@@ -12,7 +12,6 @@ struct FFZlibLibrary
FF_LIBRARY_SYMBOL(inflateInit2_)
FF_LIBRARY_SYMBOL(inflate)
FF_LIBRARY_SYMBOL(inflateEnd)
FF_LIBRARY_SYMBOL(inflateGetHeader)
bool inited;
} zlibData;
@@ -32,7 +31,6 @@ const char* ffNetworkingLoadZlibLibrary(void)
FF_LIBRARY_LOAD_SYMBOL_VAR_MESSAGE(zlib, zlibData, inflateInit2_)
FF_LIBRARY_LOAD_SYMBOL_VAR_MESSAGE(zlib, zlibData, inflate)
FF_LIBRARY_LOAD_SYMBOL_VAR_MESSAGE(zlib, zlibData, inflateEnd)
FF_LIBRARY_LOAD_SYMBOL_VAR_MESSAGE(zlib, zlibData, inflateGetHeader)
zlib = NULL; // don't auto dlclose
}
return zlibData.ffinflateEnd == NULL ? "Failed to load libz" : NULL;
@@ -69,6 +67,8 @@ static uint32_t guessGzipOutputSize(const void* data, uint32_t dataSize)
// Decompress gzip content
bool ffNetworkingDecompressGzip(FFstrbuf* buffer, char* headerEnd)
{
assert(headerEnd != NULL && *headerEnd == '\r');
// Calculate header size
uint32_t headerSize = (uint32_t) (headerEnd - buffer->chars);
@@ -86,15 +86,15 @@ bool ffNetworkingDecompressGzip(FFstrbuf* buffer, char* headerEnd)
const char* bodyStart = headerEnd + 4; // Skip delimiter
// Calculate compressed content size
uint32_t compressedSize = buffer->length - headerSize - 4;
if (compressedSize <= 0) {
if (buffer->length <= headerSize + 4) {
// No content to decompress
FF_DEBUG("Compressed content size is 0, skipping decompression");
return true;
}
// Calculate compressed content size
uint32_t compressedSize = buffer->length - headerSize - 4;
// Check if content is actually in gzip format (gzip header magic is 0x1f 0x8b)
if (compressedSize < 2 || (uint8_t)bodyStart[0] != 0x1f || (uint8_t)bodyStart[1] != 0x8b) {
FF_DEBUG("Content is not valid gzip format, skipping decompression");
@@ -137,7 +137,7 @@ bool ffNetworkingDecompressGzip(FFstrbuf* buffer, char* headerEnd)
// Save already decompressed data amount
uint32_t alreadyDecompressed = (uint32_t)(availableOut - zs.avail_out);
decompressedBuffer.length += alreadyDecompressed;
decompressedBuffer.chars[decompressedBuffer.length] = '\0'; // Ensure null-terminated string
decompressedBuffer.chars[decompressedBuffer.length] = '\0';
ffStrbufEnsureFree(&decompressedBuffer, decompressedBuffer.length / 2);
@@ -155,7 +155,7 @@ bool ffNetworkingDecompressGzip(FFstrbuf* buffer, char* headerEnd)
// Calculate decompressed size
uint32_t decompressedSize = (uint32_t)(availableOut - zs.avail_out);
decompressedBuffer.length += decompressedSize;
decompressedBuffer.chars[decompressedSize] = '\0'; // Ensure null-terminated string
decompressedBuffer.chars[decompressedBuffer.length] = '\0';
FF_DEBUG("Successfully decompressed %u bytes compressed data to %u bytes", compressedSize, decompressedBuffer.length);
// Modify Content-Length header and remove Content-Encoding header
@@ -171,7 +171,7 @@ bool ffNetworkingDecompressGzip(FFstrbuf* buffer, char* headerEnd)
}
else if (ffStrStartsWithIgnCase(line, "Content-Length:"))
{
ffStrbufAppendF(&newBuffer, "Content-Length: %u\r\n", decompressedSize);
ffStrbufAppendF(&newBuffer, "Content-Length: %u\r\n", decompressedBuffer.length);
continue;
}
else if (line[0] == '\r')
@@ -181,7 +181,7 @@ bool ffNetworkingDecompressGzip(FFstrbuf* buffer, char* headerEnd)
break;
}
ffStrbufAppendS(&newBuffer, line);
ffStrbufAppendS(&newBuffer, line); // Including the trailing \r
ffStrbufAppendC(&newBuffer, '\n');
}
+39 -31
View File
@@ -34,11 +34,7 @@ static const char* tryNonThreadingFastPath(FFNetworkingState* state)
#ifndef __APPLE__ // On macOS, TCP_FASTOPEN doesn't seem to be needed
// Set TCP Fast Open
#if __linux__ || __GNU__
int flag = 5; // the queue length of pending packets
#else
int flag = 1; // enable TCP Fast Open
#endif
int flag = 1;
if (setsockopt(state->sockfd, IPPROTO_TCP,
#ifdef __APPLE__
// https://github.com/rust-lang/libc/pull/3135
@@ -186,15 +182,12 @@ static const char* initNetworkingState(FFNetworkingState* state, const char* hos
FF_DEBUG("Initializing network connection state: host=%s, path=%s", host, path);
// Initialize command and host information
ffStrbufInitA(&state->command, 64);
ffStrbufInitA(&state->command, 128);
ffStrbufAppendS(&state->command, "GET ");
ffStrbufAppendS(&state->command, path);
ffStrbufAppendS(&state->command, " HTTP/1.0\nHost: ");
ffStrbufAppendS(&state->command, " HTTP/1.0\r\nHost: ");
ffStrbufAppendS(&state->command, host);
ffStrbufAppendS(&state->command, "\r\n");
// Add extra optimized HTTP headers
ffStrbufAppendS(&state->command, "Connection: close\r\n"); // Explicitly tell the server we don't need to keep the connection
ffStrbufAppendS(&state->command, "\r\nConnection: close\r\n"); // Explicitly tell the server we don't need to keep the connection
// If compression needs to be enabled
if (state->compression) {
@@ -221,9 +214,10 @@ static const char* initNetworkingState(FFNetworkingState* state, const char* hos
FF_DEBUG("Resolving address: %s (%s)", host, state->ipv6 ? "IPv6" : "IPv4");
// Use AI_NUMERICSERV flag to indicate the service is a numeric port, reducing parsing time
if(getaddrinfo(host, "80", &hints, &state->addr) != 0)
int gaiRes = getaddrinfo(host, "80", &hints, &state->addr);
if(gaiRes != 0)
{
FF_DEBUG("getaddrinfo() failed: %s", gai_strerror(errno));
FF_DEBUG("getaddrinfo() failed: %s (res=%d)", gai_strerror(gaiRes), gaiRes);
ret = "getaddrinfo() failed";
goto error;
}
@@ -360,6 +354,7 @@ const char* ffNetworkingSendHttpRequest(FFNetworkingState* state, const char* ho
const char* ffNetworkingRecvHttpResponse(FFNetworkingState* state, FFstrbuf* buffer)
{
assert(buffer->allocated > 0);
FF_DEBUG("Preparing to receive HTTP response");
uint32_t timeout = state->timeout;
@@ -390,15 +385,25 @@ const char* ffNetworkingRecvHttpResponse(FFNetworkingState* state, FFstrbuf* buf
// poll for the socket to be readable.
// Because of the non-blocking connectx() call, the connection might not be established yet
FF_DEBUG("Using poll() to check if socket is readable");
if (poll(&(struct pollfd) {
.fd = state->sockfd,
.events = POLLIN
}, 1, timeout > 0 ? (int) timeout : -1) == -1)
{
FF_DEBUG("poll() failed: %s (errno=%d)", strerror(errno), errno);
close(state->sockfd);
state->sockfd = -1;
return "poll() failed";
int pollRes = poll(&(struct pollfd) {
.fd = state->sockfd,
.events = POLLIN
}, 1, timeout > 0 ? (int) timeout : -1);
if (pollRes == 0)
{
FF_DEBUG("poll() timed out after %u ms", timeout);
close(state->sockfd);
state->sockfd = -1;
return "poll() timeout";
}
else if (pollRes == -1)
{
FF_DEBUG("poll() failed: %s (errno=%d)", strerror(errno), errno);
close(state->sockfd);
state->sockfd = -1;
return "poll() failed";
}
}
FF_DEBUG("Socket is readable, proceeding to receive data");
#else
@@ -415,12 +420,14 @@ const char* ffNetworkingRecvHttpResponse(FFNetworkingState* state, FFstrbuf* buf
FF_DEBUG("Starting data reception");
FF_MAYBE_UNUSED int recvCount = 0;
uint32_t contentLength = 0;
char* headerEnd = NULL;
uint32_t headerEnd = 0;
do {
FF_DEBUG("Data reception loop #%d, current buffer size: %u, available space: %u",
++recvCount, buffer->length, ffStrbufGetFree(buffer));
// We set `Connection: close`, so the server will close the connection when done.
// Thus we can use MSG_WAITALL to wait until the buffer is full or the connection is closed.
ssize_t received = recv(state->sockfd, buffer->chars + buffer->length, ffStrbufGetFree(buffer), MSG_WAITALL);
if (received <= 0) {
@@ -438,10 +445,11 @@ const char* ffNetworkingRecvHttpResponse(FFNetworkingState* state, FFstrbuf* buf
FF_DEBUG("Successfully received %zd bytes of data, total: %u bytes", received, buffer->length);
// Check if HTTP header end marker is found
if (headerEnd == NULL) {
headerEnd = memmem(buffer->chars, buffer->length, "\r\n\r\n", 4);
if (headerEnd != NULL) {
FF_DEBUG("Found HTTP header end marker, position: %ld", (long)(headerEnd - buffer->chars));
if (headerEnd == 0) {
char* pHeaderEnd = memmem(buffer->chars, buffer->length, "\r\n\r\n", 4);
if (pHeaderEnd) {
headerEnd = (uint32_t)(pHeaderEnd - buffer->chars);
FF_DEBUG("Found HTTP header end marker, position: %u", headerEnd);
// Check for Content-Length header to pre-allocate enough memory
const char* clHeader = strcasestr(buffer->chars, "Content-Length:");
@@ -451,7 +459,7 @@ const char* ffNetworkingRecvHttpResponse(FFNetworkingState* state, FFstrbuf* buf
FF_DEBUG("Detected Content-Length: %u, pre-allocating buffer", contentLength);
// Ensure buffer is large enough, adding header size and some margin
ffStrbufEnsureFree(buffer, contentLength + 16);
FF_DEBUG("Extended receive buffer to %u bytes", buffer->length);
FF_DEBUG("Extended receive buffer to %u bytes", buffer->allocated);
}
}
}
@@ -467,12 +475,12 @@ const char* ffNetworkingRecvHttpResponse(FFNetworkingState* state, FFstrbuf* buf
return "Empty server response received";
}
if (headerEnd == NULL) {
if (headerEnd == 0) {
FF_DEBUG("No HTTP header end marker found");
return "No HTTP header end found";
}
if (contentLength > 0 && buffer->length != contentLength + (uint32_t)(headerEnd - buffer->chars) + 4) {
FF_DEBUG("Received content length mismatches: %u != %u", buffer->length, contentLength + (uint32_t)(headerEnd - buffer->chars) + 4);
if (contentLength > 0 && buffer->length != contentLength + headerEnd + 4) {
FF_DEBUG("Received content length mismatches: %u != %u", buffer->length, contentLength + headerEnd + 4);
return "Content length mismatch";
}
@@ -487,7 +495,7 @@ const char* ffNetworkingRecvHttpResponse(FFNetworkingState* state, FFstrbuf* buf
#ifdef FF_HAVE_ZLIB
if (state->compression) {
FF_DEBUG("Content received, checking if compressed");
if (!ffNetworkingDecompressGzip(buffer, headerEnd)) {
if (!ffNetworkingDecompressGzip(buffer, buffer->chars + headerEnd)) {
FF_DEBUG("Decompression failed or invalid compression format");
return "Failed to decompress or invalid format";
} else {
+53 -22
View File
@@ -19,6 +19,7 @@ static const char* initWsaData(WSADATA* wsaData)
if(LOBYTE(wsaData->wVersion) != 2 || HIBYTE(wsaData->wVersion) != 2) {
FF_DEBUG("Invalid wsaData version found: %d.%d", LOBYTE(wsaData->wVersion), HIBYTE(wsaData->wVersion));
WSACleanup();
return "Invalid wsaData version found";
}
@@ -26,6 +27,7 @@ static const char* initWsaData(WSADATA* wsaData)
SOCKET sockfd = socket(AF_INET, SOCK_STREAM, 0);
if(sockfd == INVALID_SOCKET) {
FF_DEBUG("socket(AF_INET, SOCK_STREAM) failed");
WSACleanup();
return "socket(AF_INET, SOCK_STREAM) failed";
}
@@ -36,6 +38,8 @@ static const char* initWsaData(WSADATA* wsaData)
&ConnectEx, sizeof(ConnectEx),
&dwBytes, NULL, NULL) != 0) {
FF_DEBUG("WSAIoctl(sockfd, SIO_GET_EXTENSION_FUNCTION_POINTER) failed");
closesocket(sockfd);
WSACleanup();
return "WSAIoctl(sockfd, SIO_GET_EXTENSION_FUNCTION_POINTER) failed";
}
@@ -145,25 +149,26 @@ const char* ffNetworkingSendHttpRequest(FFNetworkingState* state, const char* ho
}
// Initialize overlapped structure for asynchronous I/O
memset(&state->overlapped, 0, sizeof(OVERLAPPED));
state->overlapped = (OVERLAPPED){
.hEvent = CreateEventW(NULL, FALSE, FALSE, NULL)
};
// Build HTTP command
FF_STRBUF_AUTO_DESTROY command = ffStrbufCreateA(64);
ffStrbufAppendS(&command, "GET ");
ffStrbufAppendS(&command, path);
ffStrbufAppendS(&command, " HTTP/1.0\nHost: ");
ffStrbufAppendS(&command, host);
ffStrbufAppendS(&command, "\r\n");
ffStrbufAppendS(&command, "Connection: close\r\n"); // Explicitly request connection closure
ffStrbufInitA(&state->command, 128);
ffStrbufAppendS(&state->command, "GET ");
ffStrbufAppendS(&state->command, path);
ffStrbufAppendS(&state->command, " HTTP/1.0\r\nHost: ");
ffStrbufAppendS(&state->command, host);
ffStrbufAppendS(&state->command, "\r\nConnection: close\r\n"); // Explicitly request connection closure
// Add compression support if enabled
if (state->compression) {
FF_DEBUG("Enabling HTTP content compression");
ffStrbufAppendS(&command, "Accept-Encoding: gzip\r\n");
ffStrbufAppendS(&state->command, "Accept-Encoding: gzip\r\n");
}
ffStrbufAppendS(&command, headers);
ffStrbufAppendS(&command, "\r\n");
ffStrbufAppendS(&state->command, headers);
ffStrbufAppendS(&state->command, "\r\n");
#ifdef TCP_FASTOPEN
if (state->tfo)
@@ -182,10 +187,10 @@ const char* ffNetworkingSendHttpRequest(FFNetworkingState* state, const char* ho
}
#endif
FF_DEBUG("Using ConnectEx to send %u bytes of data", command.length);
FF_DEBUG("Using ConnectEx to send %u bytes of data", state->command.length);
DWORD sent = 0;
BOOL result = ConnectEx(state->sockfd, addr->ai_addr, (int)addr->ai_addrlen,
command.chars, command.length, &sent, &state->overlapped);
state->command.chars, state->command.length, &sent, &state->overlapped);
freeaddrinfo(addr);
addr = NULL;
@@ -195,8 +200,10 @@ const char* ffNetworkingSendHttpRequest(FFNetworkingState* state, const char* ho
if (WSAGetLastError() != WSA_IO_PENDING)
{
FF_DEBUG("ConnectEx() failed: %s", ffDebugWin32Error((DWORD) WSAGetLastError()));
CloseHandle(state->overlapped.hEvent);
closesocket(state->sockfd);
state->sockfd = INVALID_SOCKET;
ffStrbufDestroy(&state->command);
return "ConnectEx() failed";
}
else
@@ -215,6 +222,7 @@ const char* ffNetworkingSendHttpRequest(FFNetworkingState* state, const char* ho
const char* ffNetworkingRecvHttpResponse(FFNetworkingState* state, FFstrbuf* buffer)
{
assert(buffer->allocated > 0);
FF_DEBUG("Preparing to receive HTTP response");
if (state->sockfd == INVALID_SOCKET)
@@ -227,11 +235,13 @@ const char* ffNetworkingRecvHttpResponse(FFNetworkingState* state, FFstrbuf* buf
if (timeout > 0)
{
FF_DEBUG("WaitForSingleObject with timeout: %u ms", timeout);
if (WaitForSingleObject((HANDLE)state->sockfd, timeout) != WAIT_OBJECT_0)
if (WaitForSingleObject(state->overlapped.hEvent, timeout) != WAIT_OBJECT_0)
{
FF_DEBUG("WaitForSingleObject failed or timed out");
CancelIo((HANDLE) state->sockfd);
CloseHandle(state->overlapped.hEvent);
closesocket(state->sockfd);
ffStrbufDestroy(&state->command);
return "WaitForSingleObject(state->sockfd) failed or timeout";
}
}
@@ -241,8 +251,20 @@ const char* ffNetworkingRecvHttpResponse(FFNetworkingState* state, FFstrbuf* buf
{
FF_DEBUG("WSAGetOverlappedResult failed: %s", ffDebugWin32Error((DWORD) WSAGetLastError()));
closesocket(state->sockfd);
CloseHandle(state->overlapped.hEvent);
ffStrbufDestroy(&state->command);
return "WSAGetOverlappedResult() failed";
}
FF_DEBUG("WSAGetOverlappedResult succeeded, %u bytes sent", (unsigned) transfer);
ffStrbufDestroy(&state->command);
CloseHandle(state->overlapped.hEvent);
state->overlapped.hEvent = NULL;
if (setsockopt(state->sockfd, SOL_SOCKET, SO_UPDATE_CONNECT_CONTEXT, NULL, 0) != 0)
{
FF_DEBUG("Failed to update connect context: %s", ffDebugWin32Error((DWORD) WSAGetLastError()));
// Not a critical error, continue anyway
}
if(timeout > 0)
{
@@ -252,12 +274,16 @@ const char* ffNetworkingRecvHttpResponse(FFNetworkingState* state, FFstrbuf* buf
// Set larger receive buffer for better performance
int rcvbuf = 65536; // 64KB
setsockopt(state->sockfd, SOL_SOCKET, SO_RCVBUF, (const char*)&rcvbuf, sizeof(rcvbuf));
if (setsockopt(state->sockfd, SOL_SOCKET, SO_RCVBUF, (const char*)&rcvbuf, sizeof(rcvbuf)))
{
FF_DEBUG("Failed to set SO_RCVBUF: %s", ffDebugWin32Error((DWORD) WSAGetLastError()));
// Not a critical error, continue anyway
}
FF_DEBUG("Starting data reception");
FF_MAYBE_UNUSED int recvCount = 0;
uint32_t contentLength = 0;
char* headerEnd = NULL;
uint32_t headerEnd = 0;
do {
FF_DEBUG("Data reception loop #%d, current buffer size: %u, available space: %u",
@@ -280,10 +306,11 @@ const char* ffNetworkingRecvHttpResponse(FFNetworkingState* state, FFstrbuf* buf
FF_DEBUG("Successfully received %zd bytes of data, total: %u bytes", received, buffer->length);
// Check if HTTP header end marker is found
if (headerEnd == NULL) {
headerEnd = strstr(buffer->chars, "\r\n\r\n");
if (headerEnd != NULL) {
FF_DEBUG("Found HTTP header end marker, position: %ld", (long)(headerEnd - buffer->chars));
if (headerEnd == 0) {
char* pHeaderEnd = strstr(buffer->chars, "\r\n\r\n");
if (pHeaderEnd) {
headerEnd = (uint32_t)(pHeaderEnd - buffer->chars);
FF_DEBUG("Found HTTP header end marker, position: %u", headerEnd);
// Check for Content-Length header to pre-allocate enough memory
const char* clHeader = strcasestr(buffer->chars, "Content-Length:");
@@ -309,10 +336,14 @@ const char* ffNetworkingRecvHttpResponse(FFNetworkingState* state, FFstrbuf* buf
return "Empty server response received";
}
if (headerEnd == NULL) {
if (headerEnd == 0) {
FF_DEBUG("No HTTP header end marker found");
return "No HTTP header end found";
}
if (contentLength > 0 && buffer->length != contentLength + headerEnd + 4) {
FF_DEBUG("Received content length mismatches: %u != %u", buffer->length, contentLength + headerEnd + 4);
return "Content length mismatch";
}
if (ffStrbufStartsWithS(buffer, "HTTP/1.0 200 OK\r\n")) {
FF_DEBUG("Received valid HTTP 200 response, content length: %u bytes, total length: %u bytes",
@@ -326,7 +357,7 @@ const char* ffNetworkingRecvHttpResponse(FFNetworkingState* state, FFstrbuf* buf
#ifdef FF_HAVE_ZLIB
if (state->compression) {
FF_DEBUG("Content received, checking if compressed");
if (!ffNetworkingDecompressGzip(buffer, headerEnd)) {
if (!ffNetworkingDecompressGzip(buffer, buffer->chars + headerEnd)) {
FF_DEBUG("Decompression failed or invalid compression format");
return "Failed to decompress or invalid format";
} else {
+1 -1
View File
@@ -15,7 +15,6 @@ typedef struct FFNetworkingState {
OVERLAPPED overlapped;
#else
int sockfd;
FFstrbuf command;
struct addrinfo* addr;
#ifdef FF_HAVE_THREADS
@@ -23,6 +22,7 @@ typedef struct FFNetworkingState {
#endif
#endif
FFstrbuf command;
uint32_t timeout;
bool ipv6;
bool compression; // if true, HTTP content compression will be enabled if supported