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
constdc_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
returntrue;
@@ -134,7 +135,7 @@ class ActiveChannel {
134135
constreplacing=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
returntrue;
@@ -192,9 +193,7 @@ class Channel {
192193
this._subscribers=undefined;
193194
this._stores=undefined;
194195
this.name=name;
195-
if(typeofname==='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
returnnewTracingChannel(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+
constlinkedChannel=channel(name);
451+
linkedChannel._index=index;
452+
dc_binding.subscribers[index]=
453+
(linkedChannel._subscribers?.length||0)+
454+
(linkedChannel._stores?.size||0);
455+
returnlinkedChannel;
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
functionsetupDiagnosticsChannel(){
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
constdc=require('diagnostics_channel');
623625
constdc_binding=internalBinding('diagnostics_channel');
624-
dc_binding.linkNativeChannel((name)=>dc.channel(name));
626+
dc_binding.linkNativeChannel((name,index)=>{
627+
constchannel=dc.channel(name);
628+
channel._index=index;
629+
dc_binding.subscribers[index]=
630+
(channel._subscribers?.length||0)+
631+
(channel._stores?.size||0);
632+
returnchannel;
633+
});
625634
}
626635

627636
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)=>{

0 commit comments

Comments
 (0)
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Add copy buttons to all
 blocks\n(function() {\n function addCopyButtons() {\n document.querySelectorAll('pre code').forEach(function(codeBlock) {\n if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;\n codeBlock.parentElement.setAttribute('data-copy-added', 'true');\n \n var btn = document.createElement('button');\n btn.textContent = 'Copy';\n btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';\n btn.onmouseover = function() { this.style.opacity = '1'; };\n btn.onmouseout = function() { this.style.opacity = '0.7'; };\n btn.onclick = function() {\n navigator.clipboard.writeText(codeBlock.textContent).then(function() {\n btn.textContent = 'Copied!';\n setTimeout(function() { btn.textContent = 'Copy'; }, 1500);\n });\n };\n codeBlock.parentElement.style.position = 'relative';\n codeBlock.parentElement.appendChild(btn);\n });\n }\n \n addCopyButtons();\n \n // Re-run on dynamic content\n var observer = new MutationObserver(addCopyButtons);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Add Copy Buttons to Code Blocks");
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
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
constdc_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
returntrue;
@@ -134,7 +135,7 @@ class ActiveChannel {
134135
constreplacing=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
returntrue;
@@ -192,9 +193,7 @@ class Channel {
192193
this._subscribers=undefined;
193194
this._stores=undefined;
194195
this.name=name;
195-
if(typeofname==='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
returnnewTracingChannel(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+
constlinkedChannel=channel(name);
451+
linkedChannel._index=index;
452+
dc_binding.subscribers[index]=
453+
(linkedChannel._subscribers?.length||0)+
454+
(linkedChannel._stores?.size||0);
455+
returnlinkedChannel;
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
functionsetupDiagnosticsChannel(){
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
constdc=require('diagnostics_channel');
623625
constdc_binding=internalBinding('diagnostics_channel');
624-
dc_binding.linkNativeChannel((name)=>dc.channel(name));
626+
dc_binding.linkNativeChannel((name,index)=>{
627+
constchannel=dc.channel(name);
628+
channel._index=index;
629+
dc_binding.subscribers[index]=
630+
(channel._subscribers?.length||0)+
631+
(channel._stores?.size||0);
632+
returnchannel;
633+
});
625634
}
626635

627636
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)=>{

0 commit comments

Comments
 (0)
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Force GitHub README to respect dark mode\n(function() {\n var style = document.createElement('style');\n style.textContent = '\n .markdown-body {\n color-scheme: dark light;\n }\n .markdown-body pre { background: #161b22 !important; }\n .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; }\n .markdown-body table th, .markdown-body table td { border-color: #30363d !important; }\n .markdown-body img { background: #0d1117; }\n .markdown-body blockquote { border-left-color: #8b949e; }\n .markdown-body hr { border-color: #30363d; }\n ';\n document.head.appendChild(style);\n})();", "GitHub Dark Mode README Fix"); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
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
constdc_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
returntrue;
@@ -134,7 +135,7 @@ class ActiveChannel {
134135
constreplacing=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
returntrue;
@@ -192,9 +193,7 @@ class Channel {
192193
this._subscribers=undefined;
193194
this._stores=undefined;
194195
this.name=name;
195-
if(typeofname==='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
returnnewTracingChannel(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+
constlinkedChannel=channel(name);
451+
linkedChannel._index=index;
452+
dc_binding.subscribers[index]=
453+
(linkedChannel._subscribers?.length||0)+
454+
(linkedChannel._stores?.size||0);
455+
returnlinkedChannel;
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
functionsetupDiagnosticsChannel(){
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
constdc=require('diagnostics_channel');
623625
constdc_binding=internalBinding('diagnostics_channel');
624-
dc_binding.linkNativeChannel((name)=>dc.channel(name));
626+
dc_binding.linkNativeChannel((name,index)=>{
627+
constchannel=dc.channel(name);
628+
channel._index=index;
629+
dc_binding.subscribers[index]=
630+
(channel._subscribers?.length||0)+
631+
(channel._stores?.size||0);
632+
returnchannel;
633+
});
625634
}
626635

627636
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)=>{

0 commit comments

Comments
 (0)
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Highlight search terms from Google/DuckDuckGo/Bing referrer\n(function() {\n var ref = document.referrer;\n var terms = [];\n \n if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) {\n var url = new URL(ref);\n var q = url.searchParams.get('q') || url.searchParams.get('p');\n if (q) {\n terms = q.split(/\\s+/).filter(function(t) { return t.length > 2; });\n }\n }\n \n if (terms.length === 0) return;\n \n var style = document.createElement('style');\n style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }';\n document.head.appendChild(style);\n \n function highlight(node) {\n if (node.nodeType === 3) { // text node\n var text = node.textContent;\n var found = false;\n terms.forEach(function(term) {\n var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\\]\\\\]/g, '\\\\') + ')', 'gi');\n if (regex.test(text)) {\n found = true;\n var frag = document.createDocumentFragment();\n var parts = text.split(regex);\n parts.forEach(function(part, i) {\n if (i % 2 === 0) {\n frag.appendChild(document.createTextNode(part));\n } else {\n var span = document.createElement('span');\n span.className = 'userscript-highlight';\n span.textContent = part;\n frag.appendChild(span);\n }\n });\n node.parentNode.replaceChild(frag, node);\n }\n });\n } else if (node.nodeType === 1 && node.childNodes) { // element\n var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT'];\n if (!skipTags.includes(node.tagName)) {\n Array.from(node.childNodes).forEach(highlight);\n }\n }\n }\n \n highlight(document.body);\n \n // Re-highlight on dynamic content\n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1 || node.nodeType === 3) highlight(node);\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Highlight Search Terms"); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
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
constdc_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
returntrue;
@@ -134,7 +135,7 @@ class ActiveChannel {
134135
constreplacing=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
returntrue;
@@ -192,9 +193,7 @@ class Channel {
192193
this._subscribers=undefined;
193194
this._stores=undefined;
194195
this.name=name;
195-
if(typeofname==='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
returnnewTracingChannel(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+
constlinkedChannel=channel(name);
451+
linkedChannel._index=index;
452+
dc_binding.subscribers[index]=
453+
(linkedChannel._subscribers?.length||0)+
454+
(linkedChannel._stores?.size||0);
455+
returnlinkedChannel;
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
functionsetupDiagnosticsChannel(){
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
constdc=require('diagnostics_channel');
623625
constdc_binding=internalBinding('diagnostics_channel');
624-
dc_binding.linkNativeChannel((name)=>dc.channel(name));
626+
dc_binding.linkNativeChannel((name,index)=>{
627+
constchannel=dc.channel(name);
628+
channel._index=index;
629+
dc_binding.subscribers[index]=
630+
(channel._subscribers?.length||0)+
631+
(channel._stores?.size||0);
632+
returnchannel;
633+
});
625634
}
626635

627636
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)=>{

0 commit comments

Comments
 (0)
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Strip utm_, fbclid, gclid, etc. from all links on page\n(function() {\n var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content',\n 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid',\n 'ref', 'ref_src', 'source', 'medium', 'campaign'];\n \n function cleanUrl(url) {\n try {\n var u = new URL(url, window.location.origin);\n var changed = false;\n trackingParams.forEach(function(p) {\n if (u.searchParams.has(p)) {\n u.searchParams.delete(p);\n changed = true;\n }\n });\n return changed ? u.toString() : url;\n } catch (e) {\n return url;\n }\n }\n \n function cleanLinks() {\n document.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n \n cleanLinks();\n \n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1) {\n if (node.tagName === 'A') cleanLinks();\n node.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Remove Tracking Parameters from Links"); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + '
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
constdc_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
returntrue;
@@ -134,7 +135,7 @@ class ActiveChannel {
134135
constreplacing=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
returntrue;
@@ -192,9 +193,7 @@ class Channel {
192193
this._subscribers=undefined;
193194
this._stores=undefined;
194195
this.name=name;
195-
if(typeofname==='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
returnnewTracingChannel(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+
constlinkedChannel=channel(name);
451+
linkedChannel._index=index;
452+
dc_binding.subscribers[index]=
453+
(linkedChannel._subscribers?.length||0)+
454+
(linkedChannel._stores?.size||0);
455+
returnlinkedChannel;
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
functionsetupDiagnosticsChannel(){
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
constdc=require('diagnostics_channel');
623625
constdc_binding=internalBinding('diagnostics_channel');
624-
dc_binding.linkNativeChannel((name)=>dc.channel(name));
626+
dc_binding.linkNativeChannel((name,index)=>{
627+
constchannel=dc.channel(name);
628+
channel._index=index;
629+
dc_binding.subscribers[index]=
630+
(channel._subscribers?.length||0)+
631+
(channel._stores?.size||0);
632+
returnchannel;
633+
});
625634
}
626635

627636
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)=>{

0 commit comments

Comments
 (0)
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Auto-enable theater mode on YouTube\n(function() {\n function tryTheater() {\n var btn = document.querySelector('button[aria-label=\"Theater mode\"], ytd-player #player button[title=\"Theater mode\"]');\n if (btn && !btn.classList.contains('activated')) {\n btn.click();\n }\n }\n \n // Try immediately\n tryTheater();\n \n // Try after navigation (SPA)\n var lastUrl = location.href;\n setInterval(function() {\n if (location.href !== lastUrl) {\n lastUrl = location.href;\n setTimeout(tryTheater, 500);\n }\n }, 1000);\n \n // Also try on player load\n var observer = new MutationObserver(tryTheater);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "YouTube Theater Mode Default"); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
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
constdc_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
returntrue;
@@ -134,7 +135,7 @@ class ActiveChannel {
134135
constreplacing=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
returntrue;
@@ -192,9 +193,7 @@ class Channel {
192193
this._subscribers=undefined;
193194
this._stores=undefined;
194195
this.name=name;
195-
if(typeofname==='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
returnnewTracingChannel(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+
constlinkedChannel=channel(name);
451+
linkedChannel._index=index;
452+
dc_binding.subscribers[index]=
453+
(linkedChannel._subscribers?.length||0)+
454+
(linkedChannel._stores?.size||0);
455+
returnlinkedChannel;
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
functionsetupDiagnosticsChannel(){
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
constdc=require('diagnostics_channel');
623625
constdc_binding=internalBinding('diagnostics_channel');
624-
dc_binding.linkNativeChannel((name)=>dc.channel(name));
626+
dc_binding.linkNativeChannel((name,index)=>{
627+
constchannel=dc.channel(name);
628+
channel._index=index;
629+
dc_binding.subscribers[index]=
630+
(channel._subscribers?.length||0)+
631+
(channel._stores?.size||0);
632+
returnchannel;
633+
});
625634
}
626635

627636
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)=>{

0 commit comments

Comments
 (0)
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Remove or un-stick sticky/fixed headers that block content\n(function() {\n function unstick() {\n document.querySelectorAll('header, nav, [role=\"banner\"], .header, .navbar, .sticky, .fixed-top, [style*=\"position: fixed\"], [style*=\"position:sticky\"]').forEach(function(el) {\n if (el.style.position === 'fixed' || el.style.position === 'sticky' || \n getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') {\n el.style.position = 'static';\n el.style.top = 'auto';\n el.style.zIndex = 'auto';\n }\n });\n }\n \n unstick();\n \n var observer = new MutationObserver(unstick);\n observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] });\n})();", "Kill Sticky Headers"); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
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
constdc_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
returntrue;
@@ -134,7 +135,7 @@ class ActiveChannel {
134135
constreplacing=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
returntrue;
@@ -192,9 +193,7 @@ class Channel {
192193
this._subscribers=undefined;
193194
this._stores=undefined;
194195
this.name=name;
195-
if(typeofname==='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
returnnewTracingChannel(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+
constlinkedChannel=channel(name);
451+
linkedChannel._index=index;
452+
dc_binding.subscribers[index]=
453+
(linkedChannel._subscribers?.length||0)+
454+
(linkedChannel._stores?.size||0);
455+
returnlinkedChannel;
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
functionsetupDiagnosticsChannel(){
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
constdc=require('diagnostics_channel');
623625
constdc_binding=internalBinding('diagnostics_channel');
624-
dc_binding.linkNativeChannel((name)=>dc.channel(name));
626+
dc_binding.linkNativeChannel((name,index)=>{
627+
constchannel=dc.channel(name);
628+
channel._index=index;
629+
dc_binding.subscribers[index]=
630+
(channel._subscribers?.length||0)+
631+
(channel._stores?.size||0);
632+
returnchannel;
633+
});
625634
}
626635

627636
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)=>{

0 commit comments

Comments
 (0)
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Universal Dark Mode - works on any site\n(function() {\n var enabled = true;\n \n function applyDarkMode() {\n if (!enabled) return;\n \n // Create style element if it doesn't exist\n var style = document.getElementById('universal-dark-mode-style');\n if (!style) {\n style = document.createElement('style');\n style.id = 'universal-dark-mode-style';\n document.head.appendChild(style);\n }\n \n // Dark mode CSS - inverts colors but preserves images/video\n style.textContent = '\n /* Invert everything except media */\n html {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #1a1a2e !important;\n }\n \n /* Restore images, videos, iframes, canvas */\n img, video, iframe, canvas, svg, picture, [style*=\"background-image\"] {\n filter: invert(1) hue-rotate(180deg) !important;\n }\n \n /* Preserve specific elements that should not be inverted */\n .no-dark-mode, .no-dark-mode *,\n [data-theme=\"light\"], [data-theme=\"light\"],\n .ace_editor, .ace_editor *,\n .CodeMirror, .CodeMirror *,\n .monaco-editor, .monaco-editor *,\n .markdown-body pre, .markdown-body pre *,\n .highlight, .highlight *,\n pre code, pre code * {\n filter: none !important;\n }\n \n /* Fix common UI elements */\n .modal, .popup, .dropdown-menu, .tooltip, .popover {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #2d2d44 !important;\n border-color: #444 !important;\n }\n \n /* Scrollbars */\n ::-webkit-scrollbar { background: #1a1a2e !important; }\n ::-webkit-scrollbar-thumb { background: #444 !important; }\n ::-webkit-scrollbar-thumb:hover { background: #555 !important; }\n \n /* Selection */\n ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ';\n }\n \n function removeDarkMode() {\n var style = document.getElementById('universal-dark-mode-style');\n if (style) style.remove();\n }\n \n // Toggle with Alt+Shift+D\n document.addEventListener('keydown', function(e) {\n if (e.altKey && e.shiftKey && e.key === 'D') {\n e.preventDefault();\n enabled = !enabled;\n if (enabled) {\n applyDarkMode();\n console.log('[Universal Dark Mode] Enabled');\n } else {\n removeDarkMode();\n console.log('[Universal Dark Mode] Disabled');\n }\n }\n });\n \n // Apply on load\n applyDarkMode();\n \n // Re-apply on dynamic content\n var observer = new MutationObserver(function(mutations) {\n if (enabled && !document.getElementById('universal-dark-mode-style')) {\n applyDarkMode();\n }\n });\n observer.observe(document.head, { childList: true });\n \n console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle');\n})();", "Universal Dark Mode"); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })();
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
constdc_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
returntrue;
@@ -134,7 +135,7 @@ class ActiveChannel {
134135
constreplacing=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
returntrue;
@@ -192,9 +193,7 @@ class Channel {
192193
this._subscribers=undefined;
193194
this._stores=undefined;
194195
this.name=name;
195-
if(typeofname==='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
returnnewTracingChannel(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+
constlinkedChannel=channel(name);
451+
linkedChannel._index=index;
452+
dc_binding.subscribers[index]=
453+
(linkedChannel._subscribers?.length||0)+
454+
(linkedChannel._stores?.size||0);
455+
returnlinkedChannel;
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
functionsetupDiagnosticsChannel(){
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
constdc=require('diagnostics_channel');
623625
constdc_binding=internalBinding('diagnostics_channel');
624-
dc_binding.linkNativeChannel((name)=>dc.channel(name));
626+
dc_binding.linkNativeChannel((name,index)=>{
627+
constchannel=dc.channel(name);
628+
channel._index=index;
629+
dc_binding.subscribers[index]=
630+
(channel._subscribers?.length||0)+
631+
(channel._stores?.size||0);
632+
returnchannel;
633+
});
625634
}
626635

627636
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)=>{

0 commit comments

Comments
 (0)