diff --git a/15.pro b/15.pro index 159f599..dec5085 100644 --- a/15.pro +++ b/15.pro @@ -30,5 +30,5 @@ HEADERS += \ src/new_design/AsyncSocketData.h INCLUDEPATH += uSockets/src src -#QMAKE_CXXFLAGS += -fsanitize=address -LIBS += -lssl -lcrypto +QMAKE_CXXFLAGS += -fsanitize=address +LIBS += -lasan -lssl -lcrypto diff --git a/src/new_design/AsyncSocket.h b/src/new_design/AsyncSocket.h index 6762e31..2be9e59 100644 --- a/src/new_design/AsyncSocket.h +++ b/src/new_design/AsyncSocket.h @@ -1,20 +1,17 @@ #ifndef ASYNCSOCKET_H #define ASYNCSOCKET_H +/* This class implements async socket memory management strategies */ + #include "StaticDispatch.h" #include "LoopData.h" #include "AsyncSocketData.h" -// todo: this is where the magic happens - namespace uWS { template struct AsyncSocket : StaticDispatch { protected: - // this will probably be different on ssl and non-ssl - static const int MAX_COPY_DISTANCE = 4096; - using SOCKET_TYPE = typename StaticDispatch::SOCKET_TYPE; using StaticDispatch::static_dispatch; @@ -59,63 +56,85 @@ public: int write(const char *src, int length, bool optionally = false, int nextLength = 0) { LoopData *loopData = getLoopData(); + std::cout << "Write called with length: " << length << ", optionally: " << optionally << std::endl; + + /* Do nothing for a null sized chunk */ if (length == 0) { + std::cout << "Write returned: 0" << std::endl; return 0; } - if (loopData->corked) { + AsyncSocketData *asyncSocketData = (AsyncSocketData *) getExt(); - // prio 1: write to cork + /* Do not write anything if we have a per-socket buffer */ + if (asyncSocketData->buffer.length()) { + if (optionally) { + std::cout << "Write returned: 0" << std::endl; + return 0; + } else { + std::cout << "Buffering at top of write!" << std::endl; + + /* At least we can reserve room for next chunk if we know it up front */ + if (nextLength) { + asyncSocketData->buffer.reserve(asyncSocketData->buffer.length() + length + nextLength); + } + + /* Buffer this chunk */ + asyncSocketData->buffer.append(src, length); + std::cout << "Write returned: " << length << std::endl; + return length; + } + } + + if (loopData->corked) { + /* We are corked */ if (LoopData::CORK_BUFFER_SIZE - loopData->corkOffset >= length) { + /* If the entire chunk fits in cork buffer */ memcpy(loopData->corkBuffer + loopData->corkOffset, src, length); loopData->corkOffset += length; } else { /* Strategy differences between SSL and non-SSL */ if constexpr(SSL) { - /* For SSL we do not want to emit small chunks, so fill the cork (assuming the cork is about 16k) */ - int remainingCork = LoopData::CORK_BUFFER_SIZE - loopData->corkOffset; + /* Cork up as much as we can, optionally does not matter here as we know it will fit in cork */ + int written = write(src, std::min(LoopData::CORK_BUFFER_SIZE - loopData->corkOffset, length), false, 0); - memcpy(loopData->corkBuffer + loopData->corkOffset, src, remainingCork); - loopData->corkOffset += remainingCork; - - uncork(src + remainingCork, length - remainingCork); + /* Optionally matters here though */ + written += uncork(src + written, length - written, optionally); + std::cout << "Write returned: " << written << std::endl; + return written; } else { /* For non-SSL we take the penalty of two syscalls */ - uncork(src, length); + int written = uncork(src, length, optionally); + std::cout << "Write returned: " << written << std::endl; + return written; } } } else { - // not corked - // we should not write here if if have buffer! - int written = 0; - - AsyncSocketData *asyncSocketData = (AsyncSocketData *) getExt(); - - if (!asyncSocketData->buffer.length()) { - written = static_dispatch(us_ssl_socket_write, us_socket_write)((SOCKET_TYPE *) this, src, length, nextLength != 0); - } - + /* We are not corked */ + int written = static_dispatch(us_ssl_socket_write, us_socket_write)((SOCKET_TYPE *) this, src, length, nextLength != 0); + /* Did we fail? */ if (written < length) { - - // bail out if we can + /* If the write was optional then just bail out */ if (optionally) { + std::cout << "Write returned: " << written << std::endl; return written; } - // shit, we are fucked - std::cout << "Buffering up per-socket" << std::endl; - AsyncSocketData *asyncSocketData = (AsyncSocketData *) getExt(); + std::cout << "Buffering at bottom of write!" << std::endl; - // at least reserve enough for next failure + /* Fall back to worst possible case (should be very rare for HTTP) */ + /* At least we can reserve room for next chunk if we know it up front */ if (nextLength) { asyncSocketData->buffer.reserve(asyncSocketData->buffer.length() + length + nextLength); } + /* Buffer this chunk */ asyncSocketData->buffer.append(src, length); } } + std::cout << "Write returned: " << length << std::endl; return length; } @@ -130,32 +149,43 @@ public: write(buf, length); } - /* Uncork this socket. It is essential to remember doing this. */ - void uncork(const char *src = nullptr, int length = 0) { + /* Uncork this socket and flush or buffer any corked and/or passed data. It is essential to remember doing this. */ + /* It does NOT count bytes written from cork buffer (they are already accounted for in the write call responsible for its corking)! */ + int uncork(const char *src = nullptr, int length = 0, bool optionally = false) { LoopData *loopData = getLoopData(); if (loopData->corked) { loopData->corked = false; if (loopData->corkOffset) { + /* Corked data is already accounted for via its write call */ write(loopData->corkBuffer, loopData->corkOffset, false, length); - write(src, length, false, 0); loopData->corkOffset = 0; } + + /* We should only return with new writes, not things written to cork already */ + return write(src, length, optionally, 0); } } /* Drain any socket-buffer while also optionally sending a chunk */ - void mergeDrain(std::string_view optionalChunk) { + int mergeDrain(std::string_view optionalChunk) { - std::cout << "mergeDrain" << std::endl; + // strategy: if we have two parts and both will fit in cork buffer then cork them and recursively send them off - // the question here is: should we recursively copy things to the cork buffer and try and send them off in one go? - // it would work in a recursively manner and since optionalChunk is optional, it would never increase the buffer size + // write any per-socket buffer and optionally more - // this function is called from onWritable and combines any buffered data with the chunk + AsyncSocketData *asyncSocketData = (AsyncSocketData *) getExt(); - // the cunk is optional meaning we do not care to buffer it up + // not handled yet + if (asyncSocketData->buffer.length()) { + std::cout << "ERROR! has socket buffer!" << std::endl; + exit(0); + } + + + /* Write the optional part */ + return write(optionalChunk.data(), optionalChunk.length(), true, 0); } void close() { diff --git a/src/new_design/HttpContext.h b/src/new_design/HttpContext.h index 5c1fa66..97d5a9e 100644 --- a/src/new_design/HttpContext.h +++ b/src/new_design/HttpContext.h @@ -84,6 +84,9 @@ private: // warning: if we are in shutdown state, resetting the timer is a security issue! static_dispatch(us_ssl_socket_timeout, us_socket_timeout)((SOCKET_TYPE *) s, HTTP_IDLE_TIMEOUT_S); + HttpResponseData *httpResponseData = (HttpResponseData *) static_dispatch(us_ssl_socket_ext, us_socket_ext)((SOCKET_TYPE *) s); + httpResponseData->offset = 0; + // route it! typename uWS::HttpContextData::UserData userData = { (HttpResponse *) s, httpRequest @@ -122,7 +125,7 @@ private: std::string_view chunk = httpResponseData->outStream(httpResponseData->offset); // send, including any buffered up - asyncSocket->mergeDrain(chunk); + httpResponseData->offset += asyncSocket->mergeDrain(chunk); } else { std::cout << "We did not have any outStream!" << std::endl; diff --git a/src/new_design/HttpResponse.h b/src/new_design/HttpResponse.h index 85ce4ba..a4e5a58 100644 --- a/src/new_design/HttpResponse.h +++ b/src/new_design/HttpResponse.h @@ -42,8 +42,9 @@ public: AsyncSocket::write("Content-Length: ", 16); AsyncSocket::writeUnsigned(chunk.length()); AsyncSocket::write("\r\n\r\n", 4); - if (AsyncSocket::write(chunk.data(), chunk.length(), true) < length) { + if (int written; (written = AsyncSocket::write(chunk.data(), chunk.length(), true)) < length) { std::cout << "HttpResponse::write failed to write everything" << std::endl; + getHttpResponseData()->offset = written; getHttpResponseData()->outStream = cb; } }