Skip to content

Commit 6d60fce

Browse files
Qardaduh95
authored andcommitted
diagnostics_channel: grow native channel storage
Every string-named JavaScript channel consumed an entry in the fixed native subscriber array. Creating more than 1,024 channels triggered a CHECK and terminated the process. Allocate slots only for native publishers and grow the aliased buffer when it fills. Refresh the JavaScript view after resizing and preserve the capacity in snapshots. Signed-off-by: Stephen Belanger <admin@stephenbelanger.com> PR-URL: #64497 Reviewed-By: Rafael Gonzaga <rafael.nunu@hotmail.com> Reviewed-By: Gerhard Stöbich <deb2001-github@yahoo.de> Reviewed-By: James M Snell <jasnell@gmail.com>
1 parent 1f1adc8 commit 6d60fce

7 files changed

Lines changed: 112 additions & 36 deletions

lib/diagnostics_channel.js

Lines changed: 16 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -31,8 +31,9 @@ const {
3131

3232
const { triggerUncaughtException } = internalBinding('errors');
3333

34+
// The subscriber buffer is replaced when native channel storage grows, so it
35+
// must always be accessed through the binding instead of cached.
3436
const dc_binding = internalBinding('diagnostics_channel');
35-
const { subscribers: subscriberCounts } = dc_binding;
3637

3738
const { WeakReference } = require('internal/util');
3839

@@ -111,7 +112,7 @@ class ActiveChannel {
111112
this._subscribers = ArrayPrototypeSlice(this._subscribers);
112113
ArrayPrototypePush(this._subscribers, subscription);
113114
channels.incRef(this.name);
114-
if (this._index !== undefined) subscriberCounts[this._index]++;
115+
if (this._index !== undefined) dc_binding.subscribers[this._index]++;
115116
}
116117

117118
unsubscribe(subscription) {
@@ -124,7 +125,7 @@ class ActiveChannel {
124125
ArrayPrototypePushApply(this._subscribers, after);
125126

126127
channels.decRef(this.name);
127-
if (this._index !== undefined) subscriberCounts[this._index]--;
128+
if (this._index !== undefined) dc_binding.subscribers[this._index]--;
128129
maybeMarkInactive(this);
129130

130131
return true;
@@ -134,7 +135,7 @@ class ActiveChannel {
134135
const replacing = this._stores.has(store);
135136
if (!replacing) {
136137
channels.incRef(this.name);
137-
if (this._index !== undefined) subscriberCounts[this._index]++;
138+
if (this._index !== undefined) dc_binding.subscribers[this._index]++;
138139
}
139140
this._stores.set(store, transform);
140141
}
@@ -147,7 +148,7 @@ class ActiveChannel {
147148
this._stores.delete(store);
148149

149150
channels.decRef(this.name);
150-
if (this._index !== undefined) subscriberCounts[this._index]--;
151+
if (this._index !== undefined) dc_binding.subscribers[this._index]--;
151152
maybeMarkInactive(this);
152153

153154
return true;
@@ -192,9 +193,7 @@ class Channel {
192193
this._subscribers = undefined;
193194
this._stores = undefined;
194195
this.name = name;
195-
if (typeof name === 'string') {
196-
this._index = dc_binding.getOrCreateChannelIndex(name);
197-
}
196+
this._index = undefined;
198197

199198
channels.set(name, this);
200199
}
@@ -446,7 +445,15 @@ function tracingChannel(nameOrChannels) {
446445
return new TracingChannel(nameOrChannels);
447446
}
448447

449-
dc_binding.linkNativeChannel((name) => channel(name));
448+
// Keep in sync with setupDiagnosticsChannel() in pre_execution.js.
449+
dc_binding.linkNativeChannel((name, index) => {
450+
const linkedChannel = channel(name);
451+
linkedChannel._index = index;
452+
dc_binding.subscribers[index] =
453+
(linkedChannel._subscribers?.length || 0) +
454+
(linkedChannel._stores?.size || 0);
455+
return linkedChannel;
456+
});
450457

451458
module.exports = {
452459
channel,

lib/internal/process/pre_execution.js

Lines changed: 10 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -619,9 +619,18 @@ function initializeClusterIPC() {
619619
function setupDiagnosticsChannel() {
620620
// Re-link native channels after snapshot deserialization since
621621
// JS references are cleared during serialization.
622+
// Keep this callback in sync with the initial registration in
623+
// lib/diagnostics_channel.js.
622624
const dc = require('diagnostics_channel');
623625
const dc_binding = internalBinding('diagnostics_channel');
624-
dc_binding.linkNativeChannel((name) => dc.channel(name));
626+
dc_binding.linkNativeChannel((name, index) => {
627+
const channel = dc.channel(name);
628+
channel._index = index;
629+
dc_binding.subscribers[index] =
630+
(channel._subscribers?.length || 0) +
631+
(channel._stores?.size || 0);
632+
return channel;
633+
});
625634
}
626635

627636
function initializePermission() {

src/node_diagnostics_channel.cc

Lines changed: 20 additions & 23 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,7 @@ using v8::Function;
1616
using v8::FunctionCallbackInfo;
1717
using v8::FunctionTemplate;
1818
using v8::HandleScope;
19+
using v8::Integer;
1920
using v8::Isolate;
2021
using v8::Local;
2122
using v8::Object;
@@ -28,8 +29,10 @@ BindingData::BindingData(Realm* realm,
2829
Local<Object> wrap,
2930
InternalFieldInfo* info)
3031
: SnapshotableObject(realm, wrap, type_int),
31-
subscribers_(
32-
realm->isolate(), kMaxChannels, MAYBE_FIELD_PTR(info, subscribers)) {
32+
subscribers_(realm->isolate(),
33+
info == nullptr ? kInitialChannelCapacity
34+
: info->subscribers_capacity,
35+
MAYBE_FIELD_PTR(info, subscribers)) {
3336
if (info == nullptr) {
3437
wrap->Set(realm->context(),
3538
FIXED_ONE_BYTE_STRING(realm->isolate(), "subscribers"),
@@ -50,25 +53,20 @@ uint32_t BindingData::GetOrCreateChannelIndex(const std::string& name) {
5053
if (it != channel_indices_.end()) {
5154
return it->second;
5255
}
53-
CHECK_LT(next_channel_index_, kMaxChannels);
56+
if (next_channel_index_ == subscribers_.Length()) {
57+
subscribers_.reserve(subscribers_.Length() * 2);
58+
object()
59+
->Set(realm()->context(),
60+
FIXED_ONE_BYTE_STRING(realm()->isolate(), "subscribers"),
61+
subscribers_.GetJSArray())
62+
.Check();
63+
subscribers_.MakeWeak();
64+
}
5465
uint32_t index = next_channel_index_++;
5566
channel_indices_.emplace(name, index);
5667
return index;
5768
}
5869

59-
void BindingData::GetOrCreateChannelIndex(
60-
const FunctionCallbackInfo<Value>& args) {
61-
Realm* realm = Realm::GetCurrent(args);
62-
BindingData* binding = realm->GetBindingData<BindingData>();
63-
CHECK_NOT_NULL(binding);
64-
65-
CHECK(args[0]->IsString());
66-
Utf8Value name(realm->isolate(), args[0]);
67-
68-
uint32_t index = binding->GetOrCreateChannelIndex(*name);
69-
args.GetReturnValue().Set(index);
70-
}
71-
7270
void BindingData::LinkNativeChannel(const FunctionCallbackInfo<Value>& args) {
7371
Realm* realm = Realm::GetCurrent(args);
7472
BindingData* binding = realm->GetBindingData<BindingData>();
@@ -85,10 +83,11 @@ void BindingData::LinkNativeChannel(const FunctionCallbackInfo<Value>& args) {
8583
Local<String> name =
8684
String::NewFromUtf8(isolate, channel_ptr->name_.c_str())
8785
.ToLocalChecked();
88-
Local<Value> argv[] = {name};
86+
Local<Value> argv[] = {
87+
name, Integer::NewFromUnsigned(isolate, channel_ptr->index_)};
8988
Local<Value> result;
9089
if (binding->link_callback_.Get(isolate)
91-
->Call(context, v8::Undefined(isolate), 1, argv)
90+
->Call(context, v8::Undefined(isolate), arraysize(argv), argv)
9291
.ToLocal(&result) &&
9392
result->IsObject()) {
9493
channel_ptr->Link(isolate, result.As<Object>());
@@ -102,6 +101,7 @@ bool BindingData::PrepareForSerialization(Local<Context> context,
102101
DCHECK_NULL(internal_field_info_);
103102
internal_field_info_ = InternalFieldInfoBase::New<InternalFieldInfo>(type());
104103
internal_field_info_->subscribers = subscribers_.Serialize(context, creator);
104+
internal_field_info_->subscribers_capacity = subscribers_.Length();
105105
link_callback_.Reset();
106106
channel_wrap_template_.Reset();
107107
channels_.clear();
@@ -130,8 +130,6 @@ void BindingData::Deserialize(Local<Context> context,
130130
void BindingData::CreatePerIsolateProperties(IsolateData* isolate_data,
131131
Local<ObjectTemplate> target) {
132132
Isolate* isolate = isolate_data->isolate();
133-
SetMethod(
134-
isolate, target, "getOrCreateChannelIndex", GetOrCreateChannelIndex);
135133
SetMethod(isolate, target, "linkNativeChannel", LinkNativeChannel);
136134
}
137135

@@ -146,7 +144,6 @@ void BindingData::CreatePerContextProperties(Local<Object> target,
146144

147145
void BindingData::RegisterExternalReferences(
148146
ExternalReferenceRegistry* registry) {
149-
registry->Register(GetOrCreateChannelIndex);
150147
registry->Register(LinkNativeChannel);
151148
}
152149

@@ -226,10 +223,10 @@ Channel* Channel::Get(Environment* env, const char* name) {
226223
HandleScope handle_scope(isolate);
227224
Local<Context> context = env->context();
228225
Local<String> js_name = String::NewFromUtf8(isolate, name).ToLocalChecked();
229-
Local<Value> argv[] = {js_name};
226+
Local<Value> argv[] = {js_name, Integer::NewFromUnsigned(isolate, index)};
230227
Local<Value> result;
231228
if (binding->link_callback_.Get(isolate)
232-
->Call(context, v8::Undefined(isolate), 1, argv)
229+
->Call(context, v8::Undefined(isolate), arraysize(argv), argv)
233230
.ToLocal(&result) &&
234231
result->IsObject()) {
235232
channel->Link(isolate, result.As<Object>());

src/node_diagnostics_channel.h

Lines changed: 2 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -20,10 +20,11 @@ class Channel;
2020

2121
class BindingData : public SnapshotableObject {
2222
public:
23-
static constexpr size_t kMaxChannels = 1024;
23+
static constexpr size_t kInitialChannelCapacity = 1024;
2424

2525
struct InternalFieldInfo : public node::InternalFieldInfoBase {
2626
AliasedBufferIndex subscribers;
27+
size_t subscribers_capacity;
2728
};
2829

2930
BindingData(Realm* realm,
@@ -48,8 +49,6 @@ class BindingData : public SnapshotableObject {
4849
v8::Global<v8::FunctionTemplate> channel_wrap_template_;
4950
std::vector<BaseObjectPtr<Channel>> channels_;
5051

51-
static void GetOrCreateChannelIndex(
52-
const v8::FunctionCallbackInfo<v8::Value>& args);
5352
static void LinkNativeChannel(
5453
const v8::FunctionCallbackInfo<v8::Value>& args);
5554

test/cctest/test_diagnostics_channel.cc

Lines changed: 44 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -257,3 +257,47 @@ TEST_F(DiagnosticsChannelTest, JSChannelVisibleFromCpp) {
257257
EXPECT_TRUE(js_has_subs->IsTrue());
258258
EXPECT_TRUE(ch->HasSubscribers());
259259
}
260+
261+
// Native channels grow the shared subscriber storage past its initial
262+
// capacity without losing the state of channels that were already linked.
263+
// Updating the first and last channels after growth also verifies that JS uses
264+
// the replacement buffer instead of a stale cached reference.
265+
TEST_F(DiagnosticsChannelTest, NativeChannelsGrowSubscriberStorage) {
266+
const v8::HandleScope handle_scope(isolate_);
267+
Argv argv;
268+
Env env{handle_scope, argv};
269+
270+
SetProcessExitHandler(*env, [&](node::Environment* env_, int exit_code) {
271+
EXPECT_EQ(exit_code, 0);
272+
node::Stop(*env);
273+
});
274+
275+
node::LoadEnvironment(
276+
*env,
277+
"globalThis.__dc = require('diagnostics_channel');"
278+
"globalThis.__firstSubscriber = () => {};"
279+
"globalThis.__dc.subscribe('test:cctest:grow:0', "
280+
" globalThis.__firstSubscriber);");
281+
282+
Channel* first = Channel::Get(*env, "test:cctest:grow:0");
283+
ASSERT_NE(first, nullptr);
284+
ASSERT_TRUE(first->HasSubscribers());
285+
286+
Channel* last = nullptr;
287+
for (size_t i = 1; i <= 1024; i++) {
288+
std::string name = "test:cctest:grow:" + std::to_string(i);
289+
last = Channel::Get(*env, name.c_str());
290+
ASSERT_NE(last, nullptr);
291+
}
292+
293+
RunJS(isolate_,
294+
"globalThis.__dc.unsubscribe('test:cctest:grow:0', "
295+
" globalThis.__firstSubscriber);");
296+
EXPECT_FALSE(first->HasSubscribers());
297+
298+
RunJS(isolate_,
299+
"globalThis.__lastSubscriber = () => {};"
300+
"globalThis.__dc.subscribe('test:cctest:grow:1024', "
301+
" globalThis.__lastSubscriber);");
302+
EXPECT_TRUE(last->HasSubscribers());
303+
}
Lines changed: 19 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,19 @@
1+
'use strict';
2+
3+
const common = require('../common');
4+
const assert = require('node:assert');
5+
const dc = require('node:diagnostics_channel');
6+
7+
let last;
8+
for (let i = 0; i < 1024 * 10 + 1; i++) {
9+
last = dc.channel(`test:many-channels:${i}`);
10+
}
11+
12+
const onMessage = common.mustCall((message, name) => {
13+
assert.strictEqual(message, 'message');
14+
assert.strictEqual(name, last.name);
15+
});
16+
17+
last.subscribe(onMessage);
18+
last.publish('message');
19+
assert.strictEqual(last.unsubscribe(onMessage), true);

test/parallel/test-diagnostics-channel-symbol-named.js

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -12,6 +12,7 @@ const symbol = Symbol('test');
1212

1313
// Individual channel objects can be created to avoid future lookups
1414
const channel = dc.channel(symbol);
15+
assert.strictEqual(Object.hasOwn(channel, '_index'), true);
1516

1617
// Expect two successful publishes later
1718
channel.subscribe(common.mustCall((message, name) => {

0 commit comments

Comments
 (0)