Skip to content
Merged
9 changes: 5 additions & 4 deletions packages/bun-uws/src/AsyncSocket.h
Original file line number Diff line number Diff line change
Expand Up @@ -183,8 +183,9 @@ struct AsyncSocket {

/* Returns the user space backpressure. */
size_t getBufferedAmount() {
/* We return the actual amount of bytes in backbuffer, including pendingRemoval */
return getAsyncSocketData()->buffer.totalLength();
/* Unsent bytes; already-written bytes sitting behind the head cursor
* are not backpressure (maxBackpressure checks, drain progress). */
return getAsyncSocketData()->buffer.length();
}

/* Returns the text representation of an IPv4 or IPv6 address */
Expand Down Expand Up @@ -253,7 +254,7 @@ struct AsyncSocket {
/* Check if we couldn't write the entire buffer */
if ((unsigned int) written < buffer_len) {
/* Remove the successfully written data from the buffer */
asyncSocketData->buffer.erase((unsigned int) written);
asyncSocketData->buffer.erase((size_t) written);

/* If we wrote less than we attempted, the socket buffer is likely full
* likely is used as an optimization hint to the compiler
Expand Down Expand Up @@ -301,7 +302,7 @@ struct AsyncSocket {
/* On failure return, otherwise continue down the function */
if ((unsigned int) written < buffer_len) {
/* Update buffering (todo: we can do better here if we keep track of what happens to this guy later on) */
asyncSocketData->buffer.erase((unsigned int) written);
asyncSocketData->buffer.erase((size_t) written);

if (optionally) {
/* Thankfully we can exit early here */
Expand Down
134 changes: 100 additions & 34 deletions packages/bun-uws/src/AsyncSocketData.h
Original file line number Diff line number Diff line change
Expand Up @@ -19,53 +19,119 @@
#ifndef UWS_ASYNCSOCKETDATA_H
#define UWS_ASYNCSOCKETDATA_H

#include <string>
#include "libusockets.h"

#include <algorithm>
#include <cstdint>
#include <cstdlib>
#include <cstring>

namespace uWS {

/* Contiguous buffer with a moving head cursor: erase() bumps head,
* append()/resize() compact into the drained head gap before growing, so
* draining never memmoves or reallocates. */
struct BackPressure {
std::string buffer;
unsigned int pendingRemoval = 0;
BackPressure(BackPressure &&other) {
buffer = std::move(other.buffer);
pendingRemoval = other.pendingRemoval;
}
BackPressure() = default;
void append(const char *data, size_t length) {
buffer.append(data, length);
BackPressure(BackPressure &&other) noexcept
: buf(other.buf), head(other.head), tail(other.tail), cap(other.cap) {
other.buf = nullptr;
other.head = other.tail = other.cap = 0;
}
void erase(unsigned int length) {
pendingRemoval += length;
/* Always erase a minimum of 1/32th the current backpressure */
if (pendingRemoval > (buffer.length() >> 5)) {
buffer.erase(0, pendingRemoval);
buffer.shrink_to_fit();
pendingRemoval = 0;
}
BackPressure(const BackPressure &) = delete;
BackPressure &operator=(const BackPressure &) = delete;
~BackPressure() { us_free(buf); }

/* Unsent bytes. data() points at length() contiguous bytes. */
size_t length() const { return tail - head; }
size_t size() const { return length(); }
const char *data() const { return buf + head; }
/* Allocation footprint for memoryCost / GC reporting. */
size_t totalLength() const { return cap; }

void append(const char *src, size_t n) {
if (!n) return;
ensureTailRoom(n);
std::memcpy(buf + tail, src, n);
tail += n;
}
size_t length() {
return buffer.length() - pendingRemoval;

void erase(size_t n) {
head += n;
if (head >= tail) {
/* Fully drained: next append writes at offset 0 with no memmove. */
head = tail = 0;
release();
}
}

void clear() {
pendingRemoval = 0;
buffer.clear();
buffer.shrink_to_fit();
head = tail = 0;
release();
}
void reserve(size_t length) {
buffer.reserve(length + pendingRemoval);
}
void resize(size_t length) {
buffer.resize(length + pendingRemoval);

/* Make room for at least n live bytes without later realloc. */
void reserve(size_t n) {
if (n > length()) ensureTailRoom(n - length());
}
const char *data() {
return buffer.data() + pendingRemoval;

/* Grow to n live bytes; caller writes into data() + old length(). */
void resize(size_t n) {
size_t live = length();
if (n > live) {
ensureTailRoom(n - live);
tail += n - live;
} else {
tail = head + n;
}
}
size_t size() {
return length();

private:
static constexpr size_t MIN_CAPACITY = 4096;

char *buf = nullptr;
size_t head = 0;
size_t tail = 0;
size_t cap = 0;

/* Ensure [tail, tail+n) is writable. Prefers compacting into the drained
* head gap over growing so steady-state producer/consumer never reallocs. */
void ensureTailRoom(size_t n) {
/* tail + n and live + n cannot wrap past this; cap * 2 wrapping is
* harmless because live + n then wins the max(). */
if (n > SIZE_MAX - tail) std::abort();
if (tail + n <= cap) return;

size_t live = tail - head;
if (head && live + n <= cap) {
std::memmove(buf, buf + head, live);
head = 0;
tail = live;
return;
}

size_t newCap = std::max(std::max(cap * 2, live + n), MIN_CAPACITY);
char *nb;
if (head == 0) {
/* mi_realloc may extend in place. */
nb = (char *) us_realloc(buf, newCap);
if (!nb) std::abort();
} else {
nb = (char *) us_malloc(newCap);
if (!nb) std::abort();
if (live) std::memcpy(nb, buf + head, live);
us_free(buf);
head = 0;
tail = live;
}
Comment thread
robobun marked this conversation as resolved.
buf = nb;
cap = newCap;
}
Comment thread
coderabbitai[bot] marked this conversation as resolved.
/* The total length, incuding pending removal */
size_t totalLength() {
return buffer.length();

void release() {
us_free(buf);
buf = nullptr;
cap = 0;
Comment thread
robobun marked this conversation as resolved.
}
};

Expand Down
3 changes: 2 additions & 1 deletion packages/bun-uws/src/WebSocket.h
Original file line number Diff line number Diff line change
Expand Up @@ -106,7 +106,8 @@ struct WebSocket : AsyncSocket<SSL> {
}

size_t memoryCost() {
return getBufferedAmount() + sizeof(WebSocket);
/* Allocation footprint for reportExtraMemoryAllocated, not unsent bytes. */
return Super::getAsyncSocketData()->buffer.totalLength() + sizeof(WebSocket);
}

/* Sending fragmented messages puts a bit of effort on the user; you must not interleave regular sends
Expand Down
Loading
Loading