Skip to content

Commit 754f073

Browse files
committed
stream: allocate stream read buffers from a slab
Read buffers for streams that emit their data to JS were allocated per read: a 64KB backing store, tracked in a map, and then - since reads rarely fill the whole buffer - reallocated to the right size and copied. Allocate read buffers from a 64KB slab instead. Reads reserve the suggested size from the slab and JS receives a view over the slab's ArrayBuffer at the read's offset, using the offset mechanism that onStreamRead already supports. Unused reservation space is rewound when a read returns less than was reserved, so small reads (e.g. TLS records) share a slab. This removes the per-read allocations, the map bookkeeping and the resize copy. Signed-off-by: Matteo Collina <hello@matteocollina.com>
1 parent 7aaf9b4 commit 754f073

6 files changed

Lines changed: 206 additions & 23 deletions

File tree

‎benchmark/net/tcp-raw-c2s.js‎

Lines changed: 6 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -23,7 +23,11 @@ function main({ dur, len, type }) {
2323
TCPConnectWrap,
2424
constants: TCPConstants,
2525
} = common.binding('tcp_wrap');
26-
const { WriteWrap } = common.binding('stream_wrap');
26+
const {
27+
WriteWrap,
28+
kReadBytesOrError,
29+
streamBaseState,
30+
} = common.binding('stream_wrap');
2731
const PORT = common.PORT;
2832

2933
const serverHandle = new TCP(TCPConstants.SERVER);
@@ -55,9 +59,7 @@ function main({ dur, len, type }) {
5559
if (!buffer)
5660
fail('read');
5761

58-
// Don't slice the buffer. The point of this is to isolate, not
59-
// simulate real traffic.
60-
bytes += buffer.byteLength;
62+
bytes += streamBaseState[kReadBytesOrError];
6163
};
6264

6365
clientHandle.readStart();

‎benchmark/net/tcp-raw-pipe.js‎

Lines changed: 10 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -23,7 +23,12 @@ function main({ dur, len, type }) {
2323
TCPConnectWrap,
2424
constants: TCPConstants,
2525
} = common.binding('tcp_wrap');
26-
const { WriteWrap } = common.binding('stream_wrap');
26+
const {
27+
WriteWrap,
28+
kReadBytesOrError,
29+
kArrayBufferOffset,
30+
streamBaseState,
31+
} = common.binding('stream_wrap');
2732
const PORT = common.PORT;
2833

2934
function fail(err, syscall) {
@@ -50,9 +55,11 @@ function main({ dur, len, type }) {
5055
if (!buffer)
5156
fail('read');
5257

58+
const nread = streamBaseState[kReadBytesOrError];
59+
const offset = streamBaseState[kArrayBufferOffset];
5360
const writeReq = new WriteWrap();
5461
writeReq.async = false;
55-
err = clientHandle.writeBuffer(writeReq, Buffer.from(buffer));
62+
err = clientHandle.writeBuffer(writeReq, Buffer.from(buffer, offset, nread));
5663

5764
if (err)
5865
fail(err, 'write');
@@ -94,7 +101,7 @@ function main({ dur, len, type }) {
94101
if (!buffer)
95102
fail('read');
96103

97-
bytes += buffer.byteLength;
104+
bytes += streamBaseState[kReadBytesOrError];
98105
};
99106

100107
connectReq.oncomplete = function(err) {

‎benchmark/net/tcp-raw-s2c.js‎

Lines changed: 6 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -23,7 +23,11 @@ function main({ dur, len, type }) {
2323
TCPConnectWrap,
2424
constants: TCPConstants,
2525
} = common.binding('tcp_wrap');
26-
const { WriteWrap } = common.binding('stream_wrap');
26+
const {
27+
WriteWrap,
28+
kReadBytesOrError,
29+
streamBaseState,
30+
} = common.binding('stream_wrap');
2731
const PORT = common.PORT;
2832

2933
const serverHandle = new TCP(TCPConstants.SERVER);
@@ -116,9 +120,7 @@ function main({ dur, len, type }) {
116120
if (!buffer)
117121
fail('read');
118122

119-
// Don't slice the buffer. The point of this is to isolate, not
120-
// simulate real traffic.
121-
bytes += buffer.byteLength;
123+
bytes += streamBaseState[kReadBytesOrError];
122124
};
123125

124126
clientHandle.readStart();

‎src/env.cc‎

Lines changed: 99 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -816,6 +816,105 @@ void Environment::recycle_managed_buffer(std::unique_ptr<BackingStore> bs) {
816816
managed_buffer_cache_ = std::move(bs);
817817
}
818818

819+
uv_buf_t StreamReadSlab::Allocate(Isolate* isolate, size_t suggested) {
820+
DCHECK_GT(suggested, 0);
821+
// Reads always get the full `suggested` size: handing out a smaller
822+
// remainder would shrink the read() buffer and fragment large reads into
823+
// more system calls, which costs more than the slab tail it saves.
824+
size_t remaining = current_.bs ? current_.size() - current_.offset : 0;
825+
if (remaining < suggested) {
826+
// Retire the current slab; it stays alive through `retired_` if reads
827+
// are still pending on it, or through JS views over it otherwise.
828+
if (current_.bs && current_.pending > 0)
829+
retired_.push_back(std::move(current_));
830+
current_ = Slab();
831+
std::unique_ptr<BackingStore> bs = ArrayBuffer::NewBackingStore(
832+
isolate,
833+
std::max(kSlabSize, suggested),
834+
BackingStoreInitializationMode::kUninitialized);
835+
current_.bs = std::move(bs);
836+
}
837+
char* base = current_.data() + current_.offset;
838+
current_.offset += suggested;
839+
current_.last_base = base;
840+
current_.last_end = current_.offset;
841+
current_.pending++;
842+
return uv_buf_init(base, suggested);
843+
}
844+
845+
StreamReadSlab::Slab* StreamReadSlab::FindSlab(const char* base) {
846+
if (current_.Contains(base)) return &current_;
847+
for (Slab& slab : retired_)
848+
if (slab.Contains(base)) return &slab;
849+
return nullptr;
850+
}
851+
852+
void StreamReadSlab::CompleteReservation(Slab* slab,
853+
const uv_buf_t& buf,
854+
size_t used) {
855+
DCHECK_GT(slab->pending, 0);
856+
slab->pending--;
857+
if (buf.base == slab->last_base && slab->offset == slab->last_end) {
858+
// This was the most recent reservation and nothing was reserved after
859+
// it: rewind the unused remainder so it can be reserved again.
860+
slab->offset = (buf.base - slab->data()) + used;
861+
slab->last_end = slab->offset;
862+
}
863+
if (slab != &current_ && slab->pending == 0) {
864+
for (auto it = retired_.begin(); it != retired_.end(); ++it) {
865+
if (&*it == slab) {
866+
retired_.erase(it);
867+
break;
868+
}
869+
}
870+
}
871+
}
872+
873+
bool StreamReadSlab::Commit(Isolate* isolate,
874+
const uv_buf_t& buf,
875+
size_t nread,
876+
Local<ArrayBuffer>* ab,
877+
size_t* offset) {
878+
Slab* slab = FindSlab(buf.base);
879+
if (slab == nullptr) return false;
880+
DCHECK_LE(nread, buf.len);
881+
882+
// A partial read would leave the rest of its reservation as waste once
883+
// the slab retires - memory that counts towards V8's external memory and
884+
// drives up GC frequency. If most of the reservation would be wasted,
885+
// give the read a right-sized copy instead (as if it had never been read
886+
// into the slab) and return its reservation in full. Reads that (mostly)
887+
// fill their reservation get a zero-copy view into the slab.
888+
if (nread < buf.len - buf.len / 4 && buf.base == slab->last_base &&
889+
slab->offset == slab->last_end) {
890+
std::unique_ptr<BackingStore> bs = ArrayBuffer::NewBackingStore(
891+
isolate, nread, BackingStoreInitializationMode::kUninitialized);
892+
memcpy(bs->Data(), buf.base, nread);
893+
*ab = ArrayBuffer::New(isolate, std::move(bs));
894+
*offset = 0;
895+
CompleteReservation(slab, buf, 0);
896+
return true;
897+
}
898+
899+
*offset = buf.base - slab->data();
900+
if (slab->ab.IsEmpty()) {
901+
*ab = ArrayBuffer::New(isolate, slab->bs);
902+
slab->ab.Reset(isolate, *ab);
903+
} else {
904+
*ab = slab->ab.Get(isolate);
905+
}
906+
CompleteReservation(slab, buf, nread);
907+
return true;
908+
}
909+
910+
bool StreamReadSlab::Release(const uv_buf_t& buf) {
911+
if (buf.base == nullptr) return true;
912+
Slab* slab = FindSlab(buf.base);
913+
if (slab == nullptr) return false;
914+
CompleteReservation(slab, buf, 0);
915+
return true;
916+
}
917+
819918
std::string Environment::GetExecPath(const std::vector<std::string>& argv) {
820919
char exec_path_buf[2 * PATH_MAX];
821920
size_t exec_path_len = sizeof(exec_path_buf);

‎src/env.h‎

Lines changed: 64 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -277,6 +277,63 @@ struct ContextInfo {
277277

278278
class EnabledDebugList;
279279

280+
// A bump allocator for stream read buffers. Reads reserve a chunk of the
281+
// current slab and, once completed, are handed to JS as a view (ArrayBuffer +
282+
// offset) over the slab, so that no per-read allocation or copy is needed.
283+
// Unused reservation space is rewound when a read returns fewer bytes than
284+
// were reserved. Slabs with reads still pending when a new slab is started
285+
// (possible when multiple reads are in flight, e.g. on Windows) are kept
286+
// alive in `retired_` until those reads complete.
287+
class StreamReadSlab {
288+
public:
289+
// Sized to match the read buffer size that libuv suggests for stream
290+
// reads: a slab typically serves a single large read (still avoiding the
291+
// copy that right-sizing the buffer would need), or many small ones.
292+
// Larger slabs amortize allocations further, but stay alive (pinned by
293+
// chunk views) long enough to be promoted to V8's old generation, where
294+
// their external memory is only reclaimed by major GCs.
295+
static constexpr size_t kSlabSize = 64 * 1024;
296+
297+
// Reserve `suggested` bytes. Starts a new slab if the current one does
298+
// not have enough space left.
299+
uv_buf_t Allocate(v8::Isolate* isolate, size_t suggested);
300+
// Commit a completed read of `nread` bytes into the buffer previously
301+
// returned by Allocate(), rewinding the unused remainder of the
302+
// reservation if possible. Returns the slab's ArrayBuffer and the offset
303+
// of `buf.base` within it. Returns false if the buffer was not allocated
304+
// from this slab (e.g. reads rerouted from another stream listener).
305+
bool Commit(v8::Isolate* isolate,
306+
const uv_buf_t& buf,
307+
size_t nread,
308+
v8::Local<v8::ArrayBuffer>* ab,
309+
size_t* offset);
310+
// Return an unused reservation (failed or empty read). Returns false if
311+
// the buffer was not allocated from this slab.
312+
bool Release(const uv_buf_t& buf);
313+
314+
private:
315+
struct Slab {
316+
std::shared_ptr<v8::BackingStore> bs;
317+
v8::Global<v8::ArrayBuffer> ab;
318+
size_t offset = 0; // Bump pointer.
319+
size_t pending = 0; // Reservations not yet committed/released.
320+
char* last_base = nullptr; // Most recent reservation...
321+
size_t last_end = 0; // ...and the bump pointer after it.
322+
323+
char* data() const { return static_cast<char*>(bs->Data()); }
324+
size_t size() const { return bs->ByteLength(); }
325+
bool Contains(const char* p) const {
326+
return bs && p >= data() && p < data() + size();
327+
}
328+
};
329+
330+
Slab* FindSlab(const char* base);
331+
void CompleteReservation(Slab* slab, const uv_buf_t& buf, size_t used);
332+
333+
Slab current_;
334+
std::vector<Slab> retired_;
335+
};
336+
280337
namespace per_process {
281338
extern std::shared_ptr<KVStore> system_environment;
282339
}
@@ -1066,6 +1123,10 @@ class Environment final : public MemoryRetainer {
10661123
// Only buffers that were not exposed externally may be recycled.
10671124
void recycle_managed_buffer(std::unique_ptr<v8::BackingStore> bs);
10681125

1126+
StreamReadSlab& stream_read_slab() {
1127+
return stream_read_slab_;
1128+
}
1129+
10691130
void AddUnmanagedFd(int fd);
10701131
void RemoveUnmanagedFd(int fd);
10711132

@@ -1283,6 +1344,9 @@ class Environment final : public MemoryRetainer {
12831344
released_allocated_buffers_;
12841345
std::unique_ptr<v8::BackingStore> managed_buffer_cache_;
12851346

1347+
// Used by EmitToJSStreamListener to allocate stream read buffers.
1348+
StreamReadSlab stream_read_slab_;
1349+
12861350
v8::CpuProfiler* cpu_profiler_ = nullptr;
12871351
std::vector<v8::ProfilerId> pending_profiles_;
12881352
};

‎src/stream_base.cc‎

Lines changed: 21 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -684,7 +684,7 @@ void StreamResource::ClearError() {
684684
uv_buf_t EmitToJSStreamListener::OnStreamAlloc(size_t suggested_size) {
685685
CHECK_NOT_NULL(stream_);
686686
Environment* env = static_cast<StreamBase*>(stream_)->stream_env();
687-
return env->allocate_managed_buffer(suggested_size);
687+
return env->stream_read_slab().Allocate(env->isolate(), suggested_size);
688688
}
689689

690690
void EmitToJSStreamListener::OnStreamRead(ssize_t nread, const uv_buf_t& buf_) {
@@ -694,25 +694,34 @@ void EmitToJSStreamListener::OnStreamRead(ssize_t nread, const uv_buf_t& buf_) {
694694
Isolate* isolate = env->isolate();
695695
HandleScope handle_scope(isolate);
696696
Context::Scope context_scope(env->context());
697-
std::unique_ptr<BackingStore> bs = env->release_managed_buffer(buf_);
698697

699698
if (nread <= 0) {
700-
env->recycle_managed_buffer(std::move(bs));
699+
if (!env->stream_read_slab().Release(buf_))
700+
env->recycle_managed_buffer(env->release_managed_buffer(buf_));
701701
if (nread < 0)
702702
stream->CallJSOnreadMethod(nread, Local<ArrayBuffer>());
703703
return;
704704
}
705705

706-
CHECK_LE(static_cast<size_t>(nread), bs->ByteLength());
707-
if (static_cast<size_t>(nread) != bs->ByteLength()) {
708-
std::unique_ptr<BackingStore> old_bs = std::move(bs);
709-
bs = ArrayBuffer::NewBackingStore(
710-
isolate, nread, BackingStoreInitializationMode::kUninitialized);
711-
memcpy(bs->Data(), old_bs->Data(), nread);
712-
env->recycle_managed_buffer(std::move(old_bs));
706+
CHECK_LE(static_cast<size_t>(nread), buf_.len);
707+
size_t offset = 0;
708+
Local<ArrayBuffer> ab;
709+
if (!env->stream_read_slab().Commit(isolate, buf_, nread, &ab, &offset)) {
710+
// The buffer was allocated through allocate_managed_buffer() by another
711+
// stream listener (e.g. StreamPipe's) whose read was rerouted here after
712+
// that listener was removed.
713+
std::unique_ptr<BackingStore> bs = env->release_managed_buffer(buf_);
714+
CHECK_LE(static_cast<size_t>(nread), bs->ByteLength());
715+
if (static_cast<size_t>(nread) != bs->ByteLength()) {
716+
std::unique_ptr<BackingStore> old_bs = std::move(bs);
717+
bs = ArrayBuffer::NewBackingStore(
718+
isolate, nread, BackingStoreInitializationMode::kUninitialized);
719+
memcpy(bs->Data(), old_bs->Data(), nread);
720+
env->recycle_managed_buffer(std::move(old_bs));
721+
}
722+
ab = ArrayBuffer::New(isolate, std::move(bs));
713723
}
714-
715-
stream->CallJSOnreadMethod(nread, ArrayBuffer::New(isolate, std::move(bs)));
724+
stream->CallJSOnreadMethod(nread, ab, offset);
716725
}
717726

718727

0 commit comments

Comments
 (0)