Sitelet https://github.com/ESP32Async/ESPAsyncWebServer/commit/dac0a61457f5f5f5faa703728b04f2a3bfcb9317
Skip to content

Commit dac0a61

Browse files
committed
Lift WS send state onto AsyncWebSocketClient; enforce one outstanding frame
AsyncWebSocketMessage is now a pure payload holder (buffer + opcode + mask), with no per-instance send-progress fields. All in-flight-frame state moves to a single, flat block of scalars on AsyncWebSocketClient (_sendingControl, _sent, _maskKey, _framePayloadLen, _frameSent, _unacked) shared uniformly between control and data frames -- there's deliberately one set of these, not one per queue, enforcing that exactly one WS frame may be outstanding (added to TCP but not yet acked) at a time. Consequences of that invariant: - A message is only popped from its queue once fully sent *and* fully acked, so queueLen()/WS_MAX_QUEUED_MESSAGES need no changes -- a message in progress is still physically sitting in the deque, exactly as before. - The next queued item (control or data) is only picked up once the in-flight frame is acked, which is also what lets control frames interleave between fragments of a data message: the priority check runs at every such boundary, not just at message boundaries. - No ack-tracking FIFO is needed: since only one frame is ever outstanding, _onAck is a single clamped subtraction instead of a per-message loop -- this also fixes a latent underflow in the old control-frame ack handling (len -= head.len() was unclamped). - _clearQueue(), AsyncWebSocketMessage::ack()/send(), and AwsMessageStatus are all gone: nothing is left to sweep once queue fronts are popped directly at the point they're known to be done. Net effect is a memory-footprint reduction, not increase: the per-message fields this removes cost more (spread across up to WS_MAX_QUEUED_MESSAGES queued messages) than the one fixed block this adds to the client.
1 parent 09e9bcf commit dac0a61

2 files changed

Lines changed: 93 additions & 203 deletions

File tree

‎src/AsyncWebSocket.cpp‎

Lines changed: 68 additions & 169 deletions
Original file line numberDiff line numberDiff line change
@@ -149,104 +149,6 @@ bool AsyncWebSocketMessageBuffer::reserve(size_t size) {
149149
return _buffer->capacity() >= size;
150150
}
151151

152-
/*
153-
* AsyncWebSocketMessage Message
154-
*/
155-
156-
AsyncWebSocketMessage::AsyncWebSocketMessage(AsyncWebSocketSharedBuffer buffer, uint8_t opcode, bool mask)
157-
: _WSbuffer{buffer}, _opcode(opcode & 0x0F), _mask{mask}, _status{_WSbuffer ? WS_MSG_SENDING : WS_MSG_ERROR} {}
158-
159-
size_t AsyncWebSocketMessage::ack(size_t len, uint32_t time) {
160-
(void)time;
161-
const size_t pending = std::min(len, _ack - _acked);
162-
_acked += pending;
163-
if (_sent >= _WSbuffer->size() && _acked >= _ack) {
164-
_status = WS_MSG_SENT;
165-
}
166-
const size_t remaining = len - pending;
167-
async_ws_log_v("ACK[%" PRIu8 "] %u/%u (acked: %u/%u) => %" PRIu8, _opcode, _sent, _WSbuffer->size(), _acked, _ack, static_cast<uint8_t>(_status));
168-
return remaining;
169-
}
170-
171-
size_t AsyncWebSocketMessage::send(AsyncClient *client) {
172-
if (!client) {
173-
async_ws_log_v("No client");
174-
return 0;
175-
}
176-
177-
if (_status != WS_MSG_SENDING) {
178-
async_ws_log_v("SEND[%" PRIu8 "] => [%" PRIu16 "] WS_MSG_SENDING != %" PRIu8, _opcode, client->remotePort(), static_cast<uint8_t>(_status));
179-
return 0;
180-
}
181-
182-
if (_sent > _WSbuffer->size()) {
183-
_status = WS_MSG_ERROR;
184-
async_ws_log_v(
185-
"SEND[%" PRIu8 "] => [%" PRIu16 "] WS_MSG_ERROR %u/%u (acked: %u/%u)", _opcode, client->remotePort(), _sent, _WSbuffer->size(), _acked, _ack
186-
);
187-
return 0;
188-
}
189-
190-
// idle between frames: pick the next frame's parameters. isControl() is
191-
// always exactly one frame, so `_ack > 0` (a prior frame fully committed)
192-
// is what actually distinguishes "done" from "haven't sent the (possibly
193-
// zero-length) first frame yet" -- `_sent == size()` alone is ambiguous
194-
// for a bare, zero-payload control frame.
195-
if (_frameSent == 0 && _framePayloadLen == 0) {
196-
if (_sent == _WSbuffer->size() && _ack > 0) {
197-
if (isControl() || _acked == _ack) {
198-
_status = WS_MSG_SENT;
199-
}
200-
async_ws_log_v(
201-
"SEND[%" PRIu8 "] => [%" PRIu16 "] WS_MSG_SENT %u/%u (acked: %u/%u)", _opcode, client->remotePort(), _sent, _WSbuffer->size(), _acked, _ack
202-
);
203-
return 0;
204-
}
205-
206-
const size_t remaining = _WSbuffer->size() - _sent;
207-
if (isControl()) {
208-
_framePayloadLen = remaining;
209-
} else {
210-
const size_t window = webSocketSendFrameWindow(client);
211-
if (!window) {
212-
async_ws_log_v("SEND[%" PRIu8 "] => [%" PRIu16 "] NO_SPACE %u", _opcode, client->remotePort(), remaining);
213-
return 0;
214-
}
215-
_framePayloadLen = std::min(remaining, window);
216-
}
217-
if (_mask) {
218-
_maskKey[0] = rand() % 0xFF; // NOLINT(runtime/threadsafe_fn)
219-
_maskKey[1] = rand() % 0xFF; // NOLINT(runtime/threadsafe_fn)
220-
_maskKey[2] = rand() % 0xFF; // NOLINT(runtime/threadsafe_fn)
221-
_maskKey[3] = rand() % 0xFF; // NOLINT(runtime/threadsafe_fn)
222-
}
223-
}
224-
225-
const bool final = isControl() || (_sent + _framePayloadLen == _WSbuffer->size());
226-
const uint8_t frameOpcode = (isControl() || _sent == 0) ? _opcode : (uint8_t)WS_CONTINUATION;
227-
228-
const size_t added = webSocketAddFrame(client, final, frameOpcode, _mask, _maskKey, (uint8_t *)_WSbuffer->data() + _sent, _framePayloadLen, _frameSent);
229-
_frameSent += added;
230-
231-
const size_t frameLen = webSocketFrameHeaderLen(_framePayloadLen, _mask) + _framePayloadLen;
232-
if (_frameSent < frameLen) {
233-
async_ws_log_v("SEND[%" PRIu8 "] => [%" PRIu16 "] PARTIAL %u/%u", _opcode, client->remotePort(), _frameSent, frameLen);
234-
return added;
235-
}
236-
237-
// frame fully committed
238-
_ack += _frameSent;
239-
_sent += _framePayloadLen;
240-
_frameSent = 0;
241-
_framePayloadLen = 0;
242-
client->send(); // idempotent flush; the bytes are already durably queued regardless of outcome
243-
244-
async_ws_log_v(
245-
"SEND[%" PRIu8 "] => [%" PRIu16 "] WS_MSG_SENDING %u/%u (acked: %u/%u)", _opcode, client->remotePort(), _sent, _WSbuffer->size(), _acked, _ack
246-
);
247-
return added;
248-
}
249-
250152
/*
251153
* Async WebSocket Client
252154
*/
@@ -312,25 +214,25 @@ AsyncWebSocketClient::~AsyncWebSocketClient() {
312214
_server->_handleEvent(this, WS_EVT_DISCONNECT, NULL, NULL, 0);
313215
}
314216

315-
void AsyncWebSocketClient::_clearQueue() {
316-
while (!_messageQueue.empty() && _messageQueue.front().finished()) {
317-
_messageQueue.pop_front();
318-
}
319-
}
320-
321217
void AsyncWebSocketClient::_onAck(size_t len, uint32_t time) {
218+
(void)time;
322219
_lastMessageTime = millis();
323220

324221
asyncsrv::unique_lock_type lock(_queue_lock);
325222

326-
async_ws_log_v("[%s][%" PRIu32 "] START ACK(%u, %" PRIu32 ") Q:%u", _server->url(), _clientId, len, time, _messageQueue.size());
223+
async_ws_log_v("[%s][%" PRIu32 "] ACK(%u) unacked:%u", _server->url(), _clientId, len, _unacked);
327224

328-
if (!_controlQueue.empty()) {
329-
auto &head = _controlQueue.front();
330-
if (head.finished()) {
331-
len -= head.len();
332-
if (_status == WS_DISCONNECTING && head.opcode() == WS_DISCONNECT) {
333-
_controlQueue.pop_front();
225+
// exactly one frame can ever be outstanding, so a single clamped
226+
// subtraction is all the ack accounting this needs
227+
const bool wasAwaitingAck = (_frameSent == 0 && _unacked > 0);
228+
_unacked -= std::min(len, _unacked);
229+
230+
if (wasAwaitingAck && _unacked == 0) {
231+
// the in-flight frame is now fully sent and fully acked
232+
if (_sendingControl && !_controlQueue.empty()) {
233+
const uint8_t opcode = _controlQueue.front().opcode();
234+
_controlQueue.pop_front();
235+
if (_status == WS_DISCONNECTING && opcode == WS_DISCONNECT) {
334236
_status = WS_DISCONNECTED;
335237
async_ws_log_v("[%s][%" PRIu32 "] ACK WS_DISCONNECTED", _server->url(), _clientId);
336238
// Capture _client before unlocking: _client->close() triggers the _onDisconnect() --> _handleDisconnect() --> ~AsyncWebSocketClient() chain,
@@ -342,23 +244,12 @@ void AsyncWebSocketClient::_onAck(size_t len, uint32_t time) {
342244
}
343245
return;
344246
}
345-
_controlQueue.pop_front();
247+
} else if (!_sendingControl && !_messageQueue.empty() && _sent >= _messageQueue.front().size()) {
248+
_messageQueue.pop_front();
249+
_sent = 0;
346250
}
347251
}
348252

349-
if (len && !_messageQueue.empty()) {
350-
for (auto &msg : _messageQueue) {
351-
len = msg.ack(len, time);
352-
if (len == 0) {
353-
break;
354-
}
355-
}
356-
}
357-
358-
_clearQueue();
359-
360-
async_ws_log_v("[%s][%" PRIu32 "] END ACK(%u, %" PRIu32 ") Q:%u", _server->url(), _clientId, len, time, _messageQueue.size());
361-
362253
_runQueue();
363254
}
364255

@@ -383,56 +274,64 @@ void AsyncWebSocketClient::_runQueue() {
383274
return;
384275
}
385276

386-
_clearQueue();
387-
388-
size_t space = webSocketSendFrameWindow(_client);
277+
if (_frameSent == 0) {
278+
if (_unacked > 0) {
279+
return; // waiting for the in-flight frame to be acked before starting the next one
280+
}
389281

390-
if (space) {
391-
// control frames have priority over message frames
392-
// we can send a control frame if:
393-
// - there is no message frame in the queue, or the first message frame is between frames (all bytes sent are acked)
394-
// - the control frame is not finished (not sent yet)
395-
// - there is enough space to send the control frame (control frames are small, at most 129 bytes, so we can assume that if there is space to send it, it can be sent in one go)
396-
if (_messageQueue.empty() || _messageQueue.front().betweenFrames()) {
397-
for (auto &ctrl : _controlQueue) {
398-
if (ctrl.finished()) {
399-
continue;
400-
}
401-
if (space > (size_t)(ctrl.len() - 1)) {
402-
async_ws_log_v("[%s][%" PRIu32 "] SEND CTRL %" PRIu8, _server->url(), _clientId, ctrl.opcode());
403-
ctrl.send(_client);
404-
space = webSocketSendFrameWindow(_client);
405-
}
406-
}
282+
// idle: pick the next target by priority (control first)
283+
if (!_controlQueue.empty()) {
284+
_sendingControl = true;
285+
} else if (!_messageQueue.empty()) {
286+
_sendingControl = false;
287+
} else {
288+
return; // nothing queued
407289
}
408290

409-
// then we can send message frames if there is space
410-
if (space) {
411-
for (auto &msg : _messageQueue) {
412-
if (msg._remainingBytesToSend()) {
413-
async_ws_log_v(
414-
"[%s][%" PRIu32 "][%" PRIu8 "] SEND %u/%u (acked: %u/%u)", _server->url(), _clientId, msg._opcode, msg._sent, msg._WSbuffer->size(), msg._acked,
415-
msg._ack
416-
);
417-
418-
// will use all the remaining space, or all the remaining bytes to send, whichever is smaller
419-
msg.send(_client);
420-
space = webSocketSendFrameWindow(_client);
421-
422-
// If we haven't finished sending this message, we must stop here to preserve WebSocket ordering.
423-
// We can only pipeline subsequent messages if the current one is fully passed to TCP buffer.
424-
if (msg._remainingBytesToSend()) {
425-
async_ws_log_v("[%s][%" PRIu32 "][%" PRIu8 "] NO_SPACE", _server->url(), _clientId, msg._opcode);
426-
break;
427-
}
428-
} else if (!space) {
429-
// not enough space for another message
430-
async_ws_log_v("[%s][%" PRIu32 "] NO_SPACE", _server->url(), _clientId);
431-
break;
432-
}
291+
AsyncWebSocketMessage &target = _sendingControl ? _controlQueue.front() : _messageQueue.front();
292+
const size_t remaining = target.size() - (_sendingControl ? 0 : _sent);
293+
if (_sendingControl) {
294+
_framePayloadLen = remaining;
295+
} else {
296+
const size_t window = webSocketSendFrameWindow(_client);
297+
if (!window) {
298+
return;
433299
}
300+
_framePayloadLen = std::min(remaining, window);
301+
}
302+
if (target.mask()) {
303+
_maskKey[0] = rand() % 0xFF; // NOLINT(runtime/threadsafe_fn)
304+
_maskKey[1] = rand() % 0xFF; // NOLINT(runtime/threadsafe_fn)
305+
_maskKey[2] = rand() % 0xFF; // NOLINT(runtime/threadsafe_fn)
306+
_maskKey[3] = rand() % 0xFF; // NOLINT(runtime/threadsafe_fn)
434307
}
435308
}
309+
310+
AsyncWebSocketMessage &target = _sendingControl ? _controlQueue.front() : _messageQueue.front();
311+
const bool final = _sendingControl || (_sent + _framePayloadLen == target.size());
312+
const uint8_t frameOpcode = (_sendingControl || _sent == 0) ? target.opcode() : (uint8_t)WS_CONTINUATION;
313+
uint8_t *payload = target.data() + (_sendingControl ? 0 : _sent);
314+
315+
async_ws_log_v(
316+
"[%s][%" PRIu32 "][%" PRIu8 "] SEND ctrl:%d %u/%u", _server->url(), _clientId, frameOpcode, _sendingControl, _frameSent,
317+
webSocketFrameHeaderLen(_framePayloadLen, target.mask()) + _framePayloadLen
318+
);
319+
320+
_frameSent += webSocketAddFrame(_client, final, frameOpcode, target.mask(), _maskKey, payload, _framePayloadLen, _frameSent);
321+
322+
const size_t frameLen = webSocketFrameHeaderLen(_framePayloadLen, target.mask()) + _framePayloadLen;
323+
if (_frameSent < frameLen) {
324+
return; // partial; resume on the next trigger
325+
}
326+
327+
// frame fully committed to TCP -- now waiting for ack
328+
_unacked = _frameSent;
329+
_frameSent = 0;
330+
if (!_sendingControl) {
331+
_sent += _framePayloadLen;
332+
}
333+
_framePayloadLen = 0;
334+
_client->send();
436335
}
437336

438337
bool AsyncWebSocketClient::queueIsFull() const {

‎src/AsyncWebSocket.h‎

Lines changed: 25 additions & 34 deletions
Original file line numberDiff line numberDiff line change
@@ -94,11 +94,6 @@ typedef enum {
9494
WS_PING,
9595
WS_PONG
9696
} AwsFrameType;
97-
typedef enum {
98-
WS_MSG_SENDING,
99-
WS_MSG_SENT,
100-
WS_MSG_ERROR
101-
} AwsMessageStatus;
10297
typedef enum {
10398
WS_EVT_CONNECT,
10499
WS_EVT_DISCONNECT,
@@ -129,49 +124,35 @@ class AsyncWebSocketMessageBuffer {
129124
}
130125
};
131126

127+
// Pure payload holder -- no send-progress state. Lives only in
128+
// AsyncWebSocketClient's _controlQueue/_messageQueue; all in-flight frame
129+
// tracking lives on the client, since only one frame may be outstanding
130+
// (added to TCP but not yet acked) across both queues at a time.
132131
class AsyncWebSocketMessage {
133-
friend AsyncWebSocketClient;
134-
135132
private:
136-
size_t _remainingBytesToSend() const {
137-
return _WSbuffer->size() - _sent;
138-
}
139-
140133
AsyncWebSocketSharedBuffer _WSbuffer;
141134
uint8_t _opcode{WS_TEXT};
142135
bool _mask{false};
143-
AwsMessageStatus _status{WS_MSG_ERROR};
144-
size_t _sent{}; // payload bytes committed across already-completed frames
145-
size_t _ack{};
146-
size_t _acked{};
147-
148-
// in-flight frame state, valid only while a frame is being built/resumed
149-
size_t _frameSent{}; // header+payload bytes committed so far for the current frame
150-
size_t _framePayloadLen{}; // payload length committed to the current frame's header
151-
uint8_t _maskKey[4]{};
152136

153137
public:
154-
AsyncWebSocketMessage(AsyncWebSocketSharedBuffer buffer, uint8_t opcode = WS_TEXT, bool mask = false);
138+
AsyncWebSocketMessage(AsyncWebSocketSharedBuffer buffer, uint8_t opcode = WS_TEXT, bool mask = false)
139+
: _WSbuffer{buffer}, _opcode(opcode & 0x0F), _mask{mask} {}
155140

156141
bool isControl() const {
157142
return _opcode >= WS_DISCONNECT;
158143
}
159-
bool finished() const {
160-
return _status != WS_MSG_SENDING;
161-
}
162-
bool betweenFrames() const {
163-
return _acked == _ack;
164-
}
165144
uint8_t opcode() const {
166145
return _opcode;
167146
}
168-
// wire length of a (necessarily unfragmented, unmasked) control frame
169-
size_t len() const {
170-
return _WSbuffer->size() + 2;
147+
bool mask() const {
148+
return _mask;
149+
}
150+
size_t size() const {
151+
return _WSbuffer->size();
152+
}
153+
uint8_t *data() const {
154+
return _WSbuffer->data();
171155
}
172-
173-
size_t ack(size_t len, uint32_t time);
174-
size_t send(AsyncClient *client);
175156
};
176157

177158
class AsyncWebSocketClient {
@@ -188,12 +169,22 @@ class AsyncWebSocketClient {
188169
std::deque<AsyncWebSocketMessage> _messageQueue;
189170
bool _closeWhenFull = false;
190171

172+
// Exactly one WS frame may be outstanding (added to TCP but not yet fully
173+
// acked) at a time, shared uniformly between control and data frames --
174+
// deliberately one set of fields, not one per queue. All are zero/false
175+
// when idle (nothing being built, sent, or awaited).
176+
bool _sendingControl{false}; // while non-idle: is the in-flight frame from _controlQueue.front() (true) or _messageQueue.front() (false)?
177+
size_t _sent{0}; // data messages only: payload bytes of _messageQueue.front() committed across already-completed frames
178+
uint8_t _maskKey[4]{}; // meaningful only while non-idle and the target message is masked
179+
size_t _framePayloadLen{0}; // payload length committed to the header of the in-flight frame
180+
size_t _frameSent{0}; // bytes of (header+in-flight payload) committed so far for the in-flight frame; 0 when idle
181+
size_t _unacked{0}; // wire bytes of the in-flight/just-completed frame not yet acked; 0 when idle
182+
191183
AwsFrameInfo _pinfo;
192184

193185
bool _queueControl(uint8_t opcode, const uint8_t *data = NULL, size_t len = 0, bool mask = false);
194186
bool _queueMessage(AsyncWebSocketSharedBuffer buffer, uint8_t opcode = WS_TEXT, bool mask = false);
195187
void _runQueue();
196-
void _clearQueue();
197188

198189
// this function is called when a text message is received, in order to copy the buffer and place a null terminator at the end of the buffer for easier handling of text messages.
199190
// Returns true on success, false on failure (e.g. memory allocation failure)

0 commit comments

Comments
 (0)