Skip to content

Commit bbd6fc5

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 44042c2 commit bbd6fc5

8 files changed

Lines changed: 118 additions & 36 deletions

‎lib/diagnostics_channel.js‎

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

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

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

3637
const{ WeakReference, kEmptyObject }=require('internal/util');
3738
const{ isPromise }=require('internal/util/types');
@@ -132,7 +133,7 @@ class ActiveChannel {
132133
this._subscribers=ArrayPrototypeSlice(this._subscribers);
133134
ArrayPrototypePush(this._subscribers,subscription);
134135
channels.incRef(this.name);
135-
if(this._index!==undefined)subscriberCounts[this._index]++;
136+
if(this._index!==undefined)dc_binding.subscribers[this._index]++;
136137
}
137138

138139
unsubscribe(subscription){
@@ -145,7 +146,7 @@ class ActiveChannel {
145146
ArrayPrototypePushApply(this._subscribers,after);
146147

147148
channels.decRef(this.name);
148-
if(this._index!==undefined)subscriberCounts[this._index]--;
149+
if(this._index!==undefined)dc_binding.subscribers[this._index]--;
149150
maybeMarkInactive(this);
150151

151152
returntrue;
@@ -155,7 +156,7 @@ class ActiveChannel {
155156
constreplacing=this._stores.has(store);
156157
if(!replacing){
157158
channels.incRef(this.name);
158-
if(this._index!==undefined)subscriberCounts[this._index]++;
159+
if(this._index!==undefined)dc_binding.subscribers[this._index]++;
159160
}
160161
this._stores.set(store,transform);
161162
}
@@ -168,7 +169,7 @@ class ActiveChannel {
168169
this._stores.delete(store);
169170

170171
channels.decRef(this.name);
171-
if(this._index!==undefined)subscriberCounts[this._index]--;
172+
if(this._index!==undefined)dc_binding.subscribers[this._index]--;
172173
maybeMarkInactive(this);
173174

174175
returntrue;
@@ -208,9 +209,7 @@ class Channel {
208209
this._subscribers=undefined;
209210
this._stores=undefined;
210211
this.name=name;
211-
if(typeofname==='string'){
212-
this._index=dc_binding.getOrCreateChannelIndex(name);
213-
}
212+
this._index=undefined;
214213

215214
channels.set(name,this);
216215
}
@@ -640,7 +639,15 @@ function tracingChannel(nameOrChannels) {
640639
returnnewTracingChannel(nameOrChannels);
641640
}
642641

643-
dc_binding.linkNativeChannel((name)=>channel(name));
642+
// Keep in sync with setupDiagnosticsChannel() in pre_execution.js.
643+
dc_binding.linkNativeChannel((name,index)=>{
644+
constlinkedChannel=channel(name);
645+
linkedChannel._index=index;
646+
dc_binding.subscribers[index]=
647+
(linkedChannel._subscribers?.length||0)+
648+
(linkedChannel._stores?.size||0);
649+
returnlinkedChannel;
650+
});
644651

645652
module.exports={
646653
channel,

‎lib/internal/process/pre_execution.js‎

Lines changed: 10 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -657,9 +657,18 @@ function initializeClusterIPC() {
657657
functionsetupDiagnosticsChannel(){
658658
// Re-link native channels after snapshot deserialization since
659659
// JS references are cleared during serialization.
660+
// Keep this callback in sync with the initial registration in
661+
// lib/diagnostics_channel.js.
660662
constdc=require('diagnostics_channel');
661663
constdc_binding=internalBinding('diagnostics_channel');
662-
dc_binding.linkNativeChannel((name)=>dc.channel(name));
664+
dc_binding.linkNativeChannel((name,index)=>{
665+
constchannel=dc.channel(name);
666+
channel._index=index;
667+
dc_binding.subscribers[index]=
668+
(channel._subscribers?.length||0)+
669+
(channel._stores?.size||0);
670+
returnchannel;
671+
});
663672
}
664673

665674
functioninitializePermission(){

‎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-
voidBindingData::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
voidBindingData::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
voidBindingData::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
voidBindingData::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
classBindingData : publicSnapshotableObject {
2222
public:
23-
staticconstexprsize_tkMaxChannels = 1024;
23+
staticconstexprsize_tkInitialChannelCapacity = 1024;
2424

2525
structInternalFieldInfo : publicnode::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-
staticvoidGetOrCreateChannelIndex(
52-
const v8::FunctionCallbackInfo<v8::Value>& args);
5352
staticvoidLinkNativeChannel(
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+
constcommon=require('../common');
4+
constassert=require('node:assert');
5+
constdc=require('node:diagnostics_channel');
6+
7+
letlast;
8+
for(leti=0;i<1024*10+1;i++){
9+
last=dc.channel(`test:many-channels:${i}`);
10+
}
11+
12+
constonMessage=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
constchannel=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)=>{

‎test/parallel/test-permission-diagnostics-channel.js‎

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -12,6 +12,12 @@ const assert = require('node:assert');
1212
constdc=require('node:diagnostics_channel');
1313
constfs=require('node:fs');
1414

15+
// JS-only channels must not consume the native subscriber storage used by the
16+
// permission audit publisher.
17+
for(leti=0;i<1024*10+1;i++){
18+
dc.channel(`test:permission:unrelated:${i}`);
19+
}
20+
1521
constmessages=[];
1622
dc.subscribe('node:permission-model:fs',(msg)=>{
1723
messages.push(msg);

0 commit comments

Comments
 (0)