| FazBrowse GitHub Viewer | Trending | | Home |
| Tools: [Download Repo ZIP] [Original HTTPS Page] |
1 parent 44042c2 commit bbd6fc5
8 files changed
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -30,8 +30,9 @@ const { | |||
| 30 | 30 | ||
| 31 | 31 | const { triggerUncaughtException } = internalBinding('errors'); | |
| 32 | 32 | ||
| 33 | + // The subscriber buffer is replaced when native channel storage grows, so it | ||
| 34 | + // must always be accessed through the binding instead of cached. | ||
| 33 | 35 | const dc_binding = internalBinding('diagnostics_channel'); | |
| 34 | - const { subscribers: subscriberCounts } = dc_binding; | ||
| 35 | 36 | ||
| 36 | 37 | const { WeakReference, kEmptyObject } = require('internal/util'); | |
| 37 | 38 | const { isPromise } = require('internal/util/types'); | |
@@ -132,7 +133,7 @@ class ActiveChannel { | |||
| 132 | 133 | this._subscribers = ArrayPrototypeSlice(this._subscribers); | |
| 133 | 134 | ArrayPrototypePush(this._subscribers, subscription); | |
| 134 | 135 | channels.incRef(this.name); | |
| 135 | - if (this._index !== undefined) subscriberCounts[this._index]++; | ||
| 136 | + if (this._index !== undefined) dc_binding.subscribers[this._index]++; | ||
| 136 | 137 | } | |
| 137 | 138 | ||
| 138 | 139 | unsubscribe(subscription) { | |
@@ -145,7 +146,7 @@ class ActiveChannel { | |||
| 145 | 146 | ArrayPrototypePushApply(this._subscribers, after); | |
| 146 | 147 | ||
| 147 | 148 | channels.decRef(this.name); | |
| 148 | - if (this._index !== undefined) subscriberCounts[this._index]--; | ||
| 149 | + if (this._index !== undefined) dc_binding.subscribers[this._index]--; | ||
| 149 | 150 | maybeMarkInactive(this); | |
| 150 | 151 | ||
| 151 | 152 | return true; | |
@@ -155,7 +156,7 @@ class ActiveChannel { | |||
| 155 | 156 | const replacing = this._stores.has(store); | |
| 156 | 157 | if (!replacing) { | |
| 157 | 158 | channels.incRef(this.name); | |
| 158 | - if (this._index !== undefined) subscriberCounts[this._index]++; | ||
| 159 | + if (this._index !== undefined) dc_binding.subscribers[this._index]++; | ||
| 159 | 160 | } | |
| 160 | 161 | this._stores.set(store, transform); | |
| 161 | 162 | } | |
@@ -168,7 +169,7 @@ class ActiveChannel { | |||
| 168 | 169 | this._stores.delete(store); | |
| 169 | 170 | ||
| 170 | 171 | channels.decRef(this.name); | |
| 171 | - if (this._index !== undefined) subscriberCounts[this._index]--; | ||
| 172 | + if (this._index !== undefined) dc_binding.subscribers[this._index]--; | ||
| 172 | 173 | maybeMarkInactive(this); | |
| 173 | 174 | ||
| 174 | 175 | return true; | |
@@ -208,9 +209,7 @@ class Channel { | |||
| 208 | 209 | this._subscribers = undefined; | |
| 209 | 210 | this._stores = undefined; | |
| 210 | 211 | this.name = name; | |
| 211 | - if (typeof name === 'string') { | ||
| 212 | - this._index = dc_binding.getOrCreateChannelIndex(name); | ||
| 213 | - } | ||
| 212 | + this._index = undefined; | ||
| 214 | 213 | ||
| 215 | 214 | channels.set(name, this); | |
| 216 | 215 | } | |
@@ -640,7 +639,15 @@ function tracingChannel(nameOrChannels) { | |||
| 640 | 639 | return new TracingChannel(nameOrChannels); | |
| 641 | 640 | } | |
| 642 | 641 | ||
| 643 | - dc_binding.linkNativeChannel((name) => channel(name)); | ||
| 642 | + // Keep in sync with setupDiagnosticsChannel() in pre_execution.js. | ||
| 643 | + dc_binding.linkNativeChannel((name, index) => { | ||
| 644 | + const linkedChannel = channel(name); | ||
| 645 | + linkedChannel._index = index; | ||
| 646 | + dc_binding.subscribers[index] = | ||
| 647 | + (linkedChannel._subscribers?.length || 0) + | ||
| 648 | + (linkedChannel._stores?.size || 0); | ||
| 649 | + return linkedChannel; | ||
| 650 | + }); | ||
| 644 | 651 | ||
| 645 | 652 | module.exports = { | |
| 646 | 653 | channel, | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -657,9 +657,18 @@ function initializeClusterIPC() { | |||
| 657 | 657 | function setupDiagnosticsChannel() { | |
| 658 | 658 | // Re-link native channels after snapshot deserialization since | |
| 659 | 659 | // JS references are cleared during serialization. | |
| 660 | + // Keep this callback in sync with the initial registration in | ||
| 661 | + // lib/diagnostics_channel.js. | ||
| 660 | 662 | const dc = require('diagnostics_channel'); | |
| 661 | 663 | const dc_binding = internalBinding('diagnostics_channel'); | |
| 662 | - dc_binding.linkNativeChannel((name) => dc.channel(name)); | ||
| 664 | + dc_binding.linkNativeChannel((name, index) => { | ||
| 665 | + const channel = dc.channel(name); | ||
| 666 | + channel._index = index; | ||
| 667 | + dc_binding.subscribers[index] = | ||
| 668 | + (channel._subscribers?.length || 0) + | ||
| 669 | + (channel._stores?.size || 0); | ||
| 670 | + return channel; | ||
| 671 | + }); | ||
| 663 | 672 | } | |
| 664 | 673 | ||
| 665 | 674 | function initializePermission() { | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -16,6 +16,7 @@ using v8::Function; | |||
| 16 | 16 | using v8::FunctionCallbackInfo; | |
| 17 | 17 | using v8::FunctionTemplate; | |
| 18 | 18 | using v8::HandleScope; | |
| 19 | + using v8::Integer; | ||
| 19 | 20 | using v8::Isolate; | |
| 20 | 21 | using v8::Local; | |
| 21 | 22 | using v8::Object; | |
@@ -28,8 +29,10 @@ BindingData::BindingData(Realm* realm, | |||
| 28 | 29 | Local<Object> wrap, | |
| 29 | 30 | InternalFieldInfo* info) | |
| 30 | 31 | : 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)) { | ||
| 33 | 36 | if (info == nullptr) { | |
| 34 | 37 | wrap->Set(realm->context(), | |
| 35 | 38 | FIXED_ONE_BYTE_STRING(realm->isolate(), "subscribers"), | |
@@ -50,25 +53,20 @@ uint32_t BindingData::GetOrCreateChannelIndex(const std::string& name) { | |||
| 50 | 53 | if (it != channel_indices_.end()) { | |
| 51 | 54 | return it->second; | |
| 52 | 55 | } | |
| 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 | + } | ||
| 54 | 65 | uint32_t index = next_channel_index_++; | |
| 55 | 66 | channel_indices_.emplace(name, index); | |
| 56 | 67 | return index; | |
| 57 | 68 | } | |
| 58 | 69 | ||
| 59 | - void BindingData::GetOrCreateChannelIndex( | ||
| 60 | - const FunctionCallbackInfo<Value>& args) { | ||
| 61 | - Realm* realm = Realm::GetCurrent(args); | ||
| 62 | - BindingData* binding = realm->GetBindingData<BindingData>(); | ||
| 63 | - CHECK_NOT_NULL(binding); | ||
| 64 | - | ||
| 65 | - CHECK(args[0]->IsString()); | ||
| 66 | - Utf8Value name(realm->isolate(), args[0]); | ||
| 67 | - | ||
| 68 | - uint32_t index = binding->GetOrCreateChannelIndex(*name); | ||
| 69 | - args.GetReturnValue().Set(index); | ||
| 70 | - } | ||
| 71 | - | ||
| 72 | 70 | void BindingData::LinkNativeChannel(const FunctionCallbackInfo<Value>& args) { | |
| 73 | 71 | Realm* realm = Realm::GetCurrent(args); | |
| 74 | 72 | BindingData* binding = realm->GetBindingData<BindingData>(); | |
@@ -85,10 +83,11 @@ void BindingData::LinkNativeChannel(const FunctionCallbackInfo<Value>& args) { | |||
| 85 | 83 | Local<String> name = | |
| 86 | 84 | String::NewFromUtf8(isolate, channel_ptr->name_.c_str()) | |
| 87 | 85 | .ToLocalChecked(); | |
| 88 | - Local<Value> argv[] = {name}; | ||
| 86 | + Local<Value> argv[] = { | ||
| 87 | + name, Integer::NewFromUnsigned(isolate, channel_ptr->index_)}; | ||
| 89 | 88 | Local<Value> result; | |
| 90 | 89 | if (binding->link_callback_.Get(isolate) | |
| 91 | - ->Call(context, v8::Undefined(isolate), 1, argv) | ||
| 90 | + ->Call(context, v8::Undefined(isolate), arraysize(argv), argv) | ||
| 92 | 91 | .ToLocal(&result) && | |
| 93 | 92 | result->IsObject()) { | |
| 94 | 93 | channel_ptr->Link(isolate, result.As<Object>()); | |
@@ -102,6 +101,7 @@ bool BindingData::PrepareForSerialization(Local<Context> context, | |||
| 102 | 101 | DCHECK_NULL(internal_field_info_); | |
| 103 | 102 | internal_field_info_ = InternalFieldInfoBase::New<InternalFieldInfo>(type()); | |
| 104 | 103 | internal_field_info_->subscribers = subscribers_.Serialize(context, creator); | |
| 104 | + internal_field_info_->subscribers_capacity = subscribers_.Length(); | ||
| 105 | 105 | link_callback_.Reset(); | |
| 106 | 106 | channel_wrap_template_.Reset(); | |
| 107 | 107 | channels_.clear(); | |
@@ -130,8 +130,6 @@ void BindingData::Deserialize(Local<Context> context, | |||
| 130 | 130 | void BindingData::CreatePerIsolateProperties(IsolateData* isolate_data, | |
| 131 | 131 | Local<ObjectTemplate> target) { | |
| 132 | 132 | Isolate* isolate = isolate_data->isolate(); | |
| 133 | - SetMethod( | ||
| 134 | - isolate, target, "getOrCreateChannelIndex", GetOrCreateChannelIndex); | ||
| 135 | 133 | SetMethod(isolate, target, "linkNativeChannel", LinkNativeChannel); | |
| 136 | 134 | } | |
| 137 | 135 | ||
@@ -146,7 +144,6 @@ void BindingData::CreatePerContextProperties(Local<Object> target, | |||
| 146 | 144 | ||
| 147 | 145 | void BindingData::RegisterExternalReferences( | |
| 148 | 146 | ExternalReferenceRegistry* registry) { | |
| 149 | - registry->Register(GetOrCreateChannelIndex); | ||
| 150 | 147 | registry->Register(LinkNativeChannel); | |
| 151 | 148 | } | |
| 152 | 149 | ||
@@ -226,10 +223,10 @@ Channel* Channel::Get(Environment* env, const char* name) { | |||
| 226 | 223 | HandleScope handle_scope(isolate); | |
| 227 | 224 | Local<Context> context = env->context(); | |
| 228 | 225 | 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)}; | ||
| 230 | 227 | Local<Value> result; | |
| 231 | 228 | if (binding->link_callback_.Get(isolate) | |
| 232 | - ->Call(context, v8::Undefined(isolate), 1, argv) | ||
| 229 | + ->Call(context, v8::Undefined(isolate), arraysize(argv), argv) | ||
| 233 | 230 | .ToLocal(&result) && | |
| 234 | 231 | result->IsObject()) { | |
| 235 | 232 | channel->Link(isolate, result.As<Object>()); | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -20,10 +20,11 @@ class Channel; | |||
| 20 | 20 | ||
| 21 | 21 | class BindingData : public SnapshotableObject { | |
| 22 | 22 | public: | |
| 23 | - static constexpr size_t kMaxChannels = 1024; | ||
| 23 | + static constexpr size_t kInitialChannelCapacity = 1024; | ||
| 24 | 24 | ||
| 25 | 25 | struct InternalFieldInfo : public node::InternalFieldInfoBase { | |
| 26 | 26 | AliasedBufferIndex subscribers; | |
| 27 | + size_t subscribers_capacity; | ||
| 27 | 28 | }; | |
| 28 | 29 | ||
| 29 | 30 | BindingData(Realm* realm, | |
@@ -48,8 +49,6 @@ class BindingData : public SnapshotableObject { | |||
| 48 | 49 | v8::Global<v8::FunctionTemplate> channel_wrap_template_; | |
| 49 | 50 | std::vector<BaseObjectPtr<Channel>> channels_; | |
| 50 | 51 | ||
| 51 | - static void GetOrCreateChannelIndex( | ||
| 52 | - const v8::FunctionCallbackInfo<v8::Value>& args); | ||
| 53 | 52 | static void LinkNativeChannel( | |
| 54 | 53 | const v8::FunctionCallbackInfo<v8::Value>& args); | |
| 55 | 54 | ||
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -257,3 +257,47 @@ TEST_F(DiagnosticsChannelTest, JSChannelVisibleFromCpp) { | |||
| 257 | 257 | EXPECT_TRUE(js_has_subs->IsTrue()); | |
| 258 | 258 | EXPECT_TRUE(ch->HasSubscribers()); | |
| 259 | 259 | } | |
| 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 | + } | ||
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -0,0 +1,19 @@ | |||
| 1 | + 'use strict'; | ||
| 2 | + | ||
| 3 | + const common = require('../common'); | ||
| 4 | + const assert = require('node:assert'); | ||
| 5 | + const dc = require('node:diagnostics_channel'); | ||
| 6 | + | ||
| 7 | + let last; | ||
| 8 | + for (let i = 0; i < 1024 * 10 + 1; i++) { | ||
| 9 | + last = dc.channel(`test:many-channels:${i}`); | ||
| 10 | + } | ||
| 11 | + | ||
| 12 | + const onMessage = common.mustCall((message, name) => { | ||
| 13 | + assert.strictEqual(message, 'message'); | ||
| 14 | + assert.strictEqual(name, last.name); | ||
| 15 | + }); | ||
| 16 | + | ||
| 17 | + last.subscribe(onMessage); | ||
| 18 | + last.publish('message'); | ||
| 19 | + assert.strictEqual(last.unsubscribe(onMessage), true); | ||
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -12,6 +12,7 @@ const symbol = Symbol('test'); | |||
| 12 | 12 | ||
| 13 | 13 | // Individual channel objects can be created to avoid future lookups | |
| 14 | 14 | const channel = dc.channel(symbol); | |
| 15 | + assert.strictEqual(Object.hasOwn(channel, '_index'), true); | ||
| 15 | 16 | ||
| 16 | 17 | // Expect two successful publishes later | |
| 17 | 18 | channel.subscribe(common.mustCall((message, name) => { | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -12,6 +12,12 @@ const assert = require('node:assert'); | |||
| 12 | 12 | const dc = require('node:diagnostics_channel'); | |
| 13 | 13 | const fs = require('node:fs'); | |
| 14 | 14 | ||
| 15 | + // JS-only channels must not consume the native subscriber storage used by the | ||
| 16 | + // permission audit publisher. | ||
| 17 | + for (let i = 0; i < 1024 * 10 + 1; i++) { | ||
| 18 | + dc.channel(`test:permission:unrelated:${i}`); | ||
| 19 | + } | ||
| 20 | + | ||
| 15 | 21 | const messages = []; | |
| 16 | 22 | dc.subscribe('node:permission-model:fs', (msg) => { | |
| 17 | 23 | messages.push(msg); | |
| Back | FazBrowse Home | New Git URL |
0 commit comments