#include "WebSocket.h" #include "Group.h" namespace uWS { template void WebSocket::send(const char *message, size_t length, OpCode opCode, void(*callback)(void *webSocket, void *data, bool cancelled, void *reserved), void *callbackData) { const int HEADER_LENGTH = WebSocketProtocol::LONG_MESSAGE_HEADER; struct TransformData { OpCode opCode; } transformData = {opCode}; struct WebSocketTransformer { static size_t estimate(const char *data, size_t length) { return length + HEADER_LENGTH; } static size_t transform(const char *src, char *dst, size_t length, TransformData transformData) { return WebSocketProtocol::formatMessage(dst, src, length, transformData.opCode, length, false); } }; sendTransformed((char *) message, length, callback, callbackData, transformData); } template typename WebSocket::PreparedMessage *WebSocket::prepareMessage(char *data, size_t length, OpCode opCode, bool compressed, void(*callback)(void *webSocket, void *data, bool cancelled, void *reserved)) { PreparedMessage *preparedMessage = new PreparedMessage; preparedMessage->buffer = new char[length + 10]; preparedMessage->length = WebSocketProtocol::formatMessage(preparedMessage->buffer, data, length, opCode, length, compressed); preparedMessage->references = 1; preparedMessage->callback = callback; return preparedMessage; } template typename WebSocket::PreparedMessage *WebSocket::prepareMessageBatch(std::vector &messages, std::vector &excludedMessages, OpCode opCode, bool compressed, void (*callback)(void *, void *, bool, void *)) { // should be sent in! size_t batchLength = 0; for (size_t i = 0; i < messages.size(); i++) { batchLength += messages[i].length(); } PreparedMessage *preparedMessage = new PreparedMessage; preparedMessage->buffer = new char[batchLength + 10 * messages.size()]; int offset = 0; for (size_t i = 0; i < messages.size(); i++) { offset += WebSocketProtocol::formatMessage(preparedMessage->buffer + offset, messages[i].data(), messages[i].length(), opCode, messages[i].length(), compressed); } preparedMessage->length = offset; preparedMessage->references = 1; preparedMessage->callback = callback; return preparedMessage; } // todo: see if this can be made a transformer instead template void WebSocket::sendPrepared(typename WebSocket::PreparedMessage *preparedMessage, void *callbackData) { preparedMessage->references++; void (*callback)(void *webSocket, void *userData, bool cancelled, void *reserved) = [](void *webSocket, void *userData, bool cancelled, void *reserved) { PreparedMessage *preparedMessage = (PreparedMessage *) userData; bool lastReference = !--preparedMessage->references; if (preparedMessage->callback) { preparedMessage->callback(webSocket, reserved, cancelled, (void *) lastReference); } if (lastReference) { delete [] preparedMessage->buffer; delete preparedMessage; } }; // candidate for fixed size pool allocator int memoryLength = sizeof(uS::SocketData::Queue::Message); int memoryIndex = getSocketData()->nodeData->getMemoryBlockIndex(memoryLength); uS::SocketData::Queue::Message *messagePtr = (uS::SocketData::Queue::Message *) getSocketData()->nodeData->getSmallMemoryBlock(memoryIndex); messagePtr->data = preparedMessage->buffer; messagePtr->length = preparedMessage->length; bool wasTransferred; if (write(messagePtr, wasTransferred)) { if (!wasTransferred) { getSocketData()->nodeData->freeSmallMemoryBlock((char *) messagePtr, memoryIndex); if (callback) { callback(*this, preparedMessage, false, callbackData); } } else { messagePtr->callback = callback; messagePtr->callbackData = preparedMessage; messagePtr->reserved = callbackData; } } else { if (callback) { callback(*this, preparedMessage, true, callbackData); } } } template void WebSocket::finalizeMessage(typename WebSocket::PreparedMessage *preparedMessage) { if (!--preparedMessage->references) { delete [] preparedMessage->buffer; delete preparedMessage; } } template void WebSocket::onData(uS::Socket s, char *data, int length) { Data *webSocketData = (Data *) s.getSocketData(); webSocketData->hasOutstandingPong = false; if (!s.isShuttingDown()) { s.cork(true); ((WebSocketProtocol *) webSocketData)->consume(data, length, s); if (!s.isClosed()) { s.cork(false); } } } template void WebSocket::terminate() { WebSocket::onEnd(*this); } template void WebSocket::close(int code, const char *message, size_t length) { static const int MAX_CLOSE_PAYLOAD = 123; length = std::min(MAX_CLOSE_PAYLOAD, length); getGroup(*this)->removeWebSocket(*this); getGroup(*this)->disconnectionHandler(*this, code, (char *) message, length); getSocketData()->shuttingDown = true; // todo: using the shared timer in the group, we can skip creating a new timer per socket // only this line and the one in Hub::connect uses the timeout feature startTimeout::onEnd>(); char closePayload[MAX_CLOSE_PAYLOAD + 2]; int closePayloadLength = WebSocketProtocol::formatClosePayload(closePayload, code, message, length); send(closePayload, closePayloadLength, OpCode::CLOSE, [](void *p, void *data, bool cancelled, void *reserved) { if (!cancelled) { Socket((uv_poll_t *) p).shutdown(); } }); } template void WebSocket::onEnd(uS::Socket s) { if (!s.isShuttingDown()) { getGroup(s)->removeWebSocket(s); getGroup(s)->disconnectionHandler(WebSocket(s), 1006, nullptr, 0); } else { s.cancelTimeout(); } Data *webSocketData = (Data *) s.getSocketData(); s.close(); while (!webSocketData->messageQueue.empty()) { uS::SocketData::Queue::Message *message = webSocketData->messageQueue.front(); if (message->callback) { message->callback(nullptr, message->callbackData, true, nullptr); } webSocketData->messageQueue.pop(); } delete webSocketData; } template struct WebSocket; template struct WebSocket; }