| FazBrowse GitHub Viewer | Trending | | Home |
| Tools: [Download Repo ZIP] [Original HTTPS Page] |
1 parent 1f1adc8 commit 6d60fce
7 files changed
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -31,8 +31,9 @@ const { | |||
| 31 | 31 | ||
| 32 | 32 | const { triggerUncaughtException } = internalBinding('errors'); | |
| 33 | 33 | ||
| 34 | + // The subscriber buffer is replaced when native channel storage grows, so it | ||
| 35 | + // must always be accessed through the binding instead of cached. | ||
| 34 | 36 | const dc_binding = internalBinding('diagnostics_channel'); | |
| 35 | - const { subscribers: subscriberCounts } = dc_binding; | ||
| 36 | 37 | ||
| 37 | 38 | const { WeakReference } = require('internal/util'); | |
| 38 | 39 | ||
@@ -111,7 +112,7 @@ class ActiveChannel { | |||
| 111 | 112 | this._subscribers = ArrayPrototypeSlice(this._subscribers); | |
| 112 | 113 | ArrayPrototypePush(this._subscribers, subscription); | |
| 113 | 114 | channels.incRef(this.name); | |
| 114 | - if (this._index !== undefined) subscriberCounts[this._index]++; | ||
| 115 | + if (this._index !== undefined) dc_binding.subscribers[this._index]++; | ||
| 115 | 116 | } | |
| 116 | 117 | ||
| 117 | 118 | unsubscribe(subscription) { | |
@@ -124,7 +125,7 @@ class ActiveChannel { | |||
| 124 | 125 | ArrayPrototypePushApply(this._subscribers, after); | |
| 125 | 126 | ||
| 126 | 127 | channels.decRef(this.name); | |
| 127 | - if (this._index !== undefined) subscriberCounts[this._index]--; | ||
| 128 | + if (this._index !== undefined) dc_binding.subscribers[this._index]--; | ||
| 128 | 129 | maybeMarkInactive(this); | |
| 129 | 130 | ||
| 130 | 131 | return true; | |
@@ -134,7 +135,7 @@ class ActiveChannel { | |||
| 134 | 135 | const replacing = this._stores.has(store); | |
| 135 | 136 | if (!replacing) { | |
| 136 | 137 | channels.incRef(this.name); | |
| 137 | - if (this._index !== undefined) subscriberCounts[this._index]++; | ||
| 138 | + if (this._index !== undefined) dc_binding.subscribers[this._index]++; | ||
| 138 | 139 | } | |
| 139 | 140 | this._stores.set(store, transform); | |
| 140 | 141 | } | |
@@ -147,7 +148,7 @@ class ActiveChannel { | |||
| 147 | 148 | this._stores.delete(store); | |
| 148 | 149 | ||
| 149 | 150 | channels.decRef(this.name); | |
| 150 | - if (this._index !== undefined) subscriberCounts[this._index]--; | ||
| 151 | + if (this._index !== undefined) dc_binding.subscribers[this._index]--; | ||
| 151 | 152 | maybeMarkInactive(this); | |
| 152 | 153 | ||
| 153 | 154 | return true; | |
@@ -192,9 +193,7 @@ class Channel { | |||
| 192 | 193 | this._subscribers = undefined; | |
| 193 | 194 | this._stores = undefined; | |
| 194 | 195 | this.name = name; | |
| 195 | - if (typeof name === 'string') { | ||
| 196 | - this._index = dc_binding.getOrCreateChannelIndex(name); | ||
| 197 | - } | ||
| 196 | + this._index = undefined; | ||
| 198 | 197 | ||
| 199 | 198 | channels.set(name, this); | |
| 200 | 199 | } | |
@@ -446,7 +445,15 @@ function tracingChannel(nameOrChannels) { | |||
| 446 | 445 | return new TracingChannel(nameOrChannels); | |
| 447 | 446 | } | |
| 448 | 447 | ||
| 449 | - dc_binding.linkNativeChannel((name) => channel(name)); | ||
| 448 | + // Keep in sync with setupDiagnosticsChannel() in pre_execution.js. | ||
| 449 | + dc_binding.linkNativeChannel((name, index) => { | ||
| 450 | + const linkedChannel = channel(name); | ||
| 451 | + linkedChannel._index = index; | ||
| 452 | + dc_binding.subscribers[index] = | ||
| 453 | + (linkedChannel._subscribers?.length || 0) + | ||
| 454 | + (linkedChannel._stores?.size || 0); | ||
| 455 | + return linkedChannel; | ||
| 456 | + }); | ||
| 450 | 457 | ||
| 451 | 458 | module.exports = { | |
| 452 | 459 | channel, | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -619,9 +619,18 @@ function initializeClusterIPC() { | |||
| 619 | 619 | function setupDiagnosticsChannel() { | |
| 620 | 620 | // Re-link native channels after snapshot deserialization since | |
| 621 | 621 | // JS references are cleared during serialization. | |
| 622 | + // Keep this callback in sync with the initial registration in | ||
| 623 | + // lib/diagnostics_channel.js. | ||
| 622 | 624 | const dc = require('diagnostics_channel'); | |
| 623 | 625 | const dc_binding = internalBinding('diagnostics_channel'); | |
| 624 | - dc_binding.linkNativeChannel((name) => dc.channel(name)); | ||
| 626 | + dc_binding.linkNativeChannel((name, index) => { | ||
| 627 | + const channel = dc.channel(name); | ||
| 628 | + channel._index = index; | ||
| 629 | + dc_binding.subscribers[index] = | ||
| 630 | + (channel._subscribers?.length || 0) + | ||
| 631 | + (channel._stores?.size || 0); | ||
| 632 | + return channel; | ||
| 633 | + }); | ||
| 625 | 634 | } | |
| 626 | 635 | ||
| 627 | 636 | 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) => { | |
| Back | FazBrowse Home | New Git URL |
0 commit comments