Multi-thread queue synchronization for AsyncWebSocketClient

This commit is contained in:
2020-12-20 02:08:39 +01:00
parent f5e439c8f8
commit 344d31c1ff
2 changed files with 302 additions and 120 deletions
+68 -6
View File
@@ -157,6 +157,63 @@ public:
virtual size_t send(AsyncClient *client) override ;
};
class PolymorphMessageContainer
{
union {
AsyncWebSocketBasicMessage basicMessage;
AsyncWebSocketMultiMessage multiMessage;
};
enum class Type : uint8_t { Basic, Multi };
const Type type;
public:
PolymorphMessageContainer() = delete;
PolymorphMessageContainer(const PolymorphMessageContainer &) = delete;
PolymorphMessageContainer &operator=(const PolymorphMessageContainer &) = delete;
PolymorphMessageContainer(const char *data, size_t len, uint8_t opcode=WS_TEXT, bool mask=false) :
type{Type::Basic}
{
new (&basicMessage) AsyncWebSocketBasicMessage{data, len, opcode, mask};
}
PolymorphMessageContainer(AsyncWebSocketMessageBuffer *buffer, uint8_t opcode=WS_TEXT, bool mask=false) :
type{Type::Multi}
{
new (&multiMessage) AsyncWebSocketMultiMessage{buffer, opcode, mask};
}
~PolymorphMessageContainer()
{
switch (type)
{
case Type::Basic: basicMessage.~AsyncWebSocketBasicMessage(); break;
case Type::Multi: multiMessage.~AsyncWebSocketMultiMessage(); break;
}
}
AsyncWebSocketMessage &get()
{
switch (type)
{
case Type::Basic: return basicMessage;
case Type::Multi: return multiMessage;
}
__builtin_unreachable();
}
const AsyncWebSocketMessage &get() const
{
switch (type)
{
case Type::Basic: return basicMessage;
case Type::Multi: return multiMessage;
}
__builtin_unreachable();
}
};
class AsyncWebSocketClient {
private:
AsyncClient *_client;
@@ -164,8 +221,10 @@ class AsyncWebSocketClient {
uint32_t _clientId;
AwsClientStatus _status;
AsyncWebLock _lock;
std::queue<AsyncWebSocketControl> _controlQueue;
LinkedList<AsyncWebSocketMessage *> _messageQueue;
std::queue<PolymorphMessageContainer> _messageQueue;
uint8_t _pstate;
AwsFrameInfo _pinfo;
@@ -173,8 +232,9 @@ class AsyncWebSocketClient {
uint32_t _lastMessageTime;
uint32_t _keepAlivePeriod;
void _queueMessage(AsyncWebSocketMessage *dataMessage);
void _queueControl(uint8_t opcode, uint8_t *data=NULL, size_t len=0, bool mask=false);
void _queueMessage(const char *data, size_t len, uint8_t opcode=WS_TEXT, bool mask=false);
void _queueMessage(AsyncWebSocketMessageBuffer *buffer, uint8_t opcode=WS_TEXT, bool mask=false);
void _runQueue();
public:
@@ -208,7 +268,8 @@ class AsyncWebSocketClient {
}
//data packets
void message(AsyncWebSocketMessage *message){ _queueMessage(message); }
void message(const char *data, size_t len, uint8_t opcode=WS_TEXT, bool mask=false) { _queueMessage(data, len, opcode, mask); }
void message(AsyncWebSocketMessageBuffer *buffer, uint8_t opcode=WS_TEXT, bool mask=false) { _queueMessage(buffer, opcode, mask); }
bool queueIsFull() const;
size_t printf(const char *format, ...) __attribute__ ((format (printf, 2, 3)));
@@ -231,7 +292,7 @@ class AsyncWebSocketClient {
void binary(const __FlashStringHelper *data, size_t len);
void binary(AsyncWebSocketMessageBuffer *buffer);
bool canSend() { return _messageQueue.length() < WS_MAX_QUEUED_MESSAGES; }
bool canSend() const;
//system callbacks (do not call)
void _onAck(size_t len, uint32_t time);
@@ -304,8 +365,9 @@ class AsyncWebSocket: public AsyncWebHandler {
void binaryAll(const __FlashStringHelper *message, size_t len);
void binaryAll(AsyncWebSocketMessageBuffer * buffer);
void message(uint32_t id, AsyncWebSocketMessage *message);
void messageAll(AsyncWebSocketMultiMessage *message);
void message(uint32_t id, const char *data, size_t len, uint8_t opcode=WS_TEXT, bool mask=false);
void message(uint32_t id, AsyncWebSocketMessageBuffer *buffer, uint8_t opcode=WS_TEXT, bool mask=false);
void messageAll(AsyncWebSocketMessageBuffer *buffer, uint8_t opcode=WS_TEXT, bool mask=false);
size_t printf(uint32_t id, const char *format, ...) __attribute__ ((format (printf, 3, 4)));
size_t printfAll(const char *format, ...) __attribute__ ((format (printf, 2, 3)));