File
Blob: src/workerd/api/node/diagnostics-channel.c++
| 1 | #include "diagnostics-channel.h" |
| 2 | |
| 3 | #include <workerd/io/io-context.h> |
| 4 | #include <workerd/io/trace.h> |
| 5 | #include <workerd/io/tracer.h> |
| 6 | #include <workerd/jsg/ser.h> |
| 7 | |
| 8 | namespace workerd::api::node { |
| 9 | |
| 10 | jsg::Value Channel::identityTransform(jsg::Lock& js, jsg::Value value) { |
| 11 | return value.addRef(js); |
| 12 | } |
| 13 | |
| 14 | Channel::Channel(jsg::Name name): name(kj::mv(name)) {} |
| 15 | |
| 16 | const jsg::Name& Channel::getName() const { |
| 17 | return name; |
| 18 | } |
| 19 | |
| 20 | bool Channel::hasSubscribers() { |
| 21 | return subscribers.size() != 0; |
| 22 | } |
| 23 | |
| 24 | void Channel::publish(jsg::Lock& js, jsg::Value message) { |
| 25 | auto snapshot = KJ_MAP(sub, subscribers) -> MessageCallback { return sub.value.addRef(js); }; |
| 26 | for (auto& cb: snapshot) { |
| 27 | cb(js, message.addRef(js), name.clone(js)); |
| 28 | } |
| 29 | |
| 30 | auto& context = IoContext::current(); |
| 31 | KJ_IF_SOME(tracer, context.getWorkerTracer()) { |
| 32 | js.tryCatch([&]() { |
| 33 | jsg::Serializer ser(js, |
| 34 | jsg::Serializer::Options{ |
| 35 | .omitHeader = false, |
| 36 | }); |
| 37 | ser.write(js, jsg::JsValue(message.getHandle(js))); |
| 38 | auto tmp = ser.release(); |
| 39 | JSG_REQUIRE(tmp.sharedArrayBuffers.size() == 0 && tmp.transferredArrayBuffers.size() == 0, |
| 40 | Error, |
| 41 | "Diagnostic events cannot be published with SharedArrayBuffer or " |
| 42 | "transferred ArrayBuffer instances"); |
| 43 | tracer.addDiagnosticChannelEvent( |
| 44 | context.getInvocationSpanContext(), context.now(), name.toString(js), kj::mv(tmp.data)); |
| 45 | }, [&](jsg::Value&& exception) { |
| 46 | jsg::JsValue jsException(exception.getHandle(js)); |
| 47 | tracer.addException(context.getInvocationSpanContext(), context.now(), kj::str("Error"), |
| 48 | kj::str("Failed to publish diagnostics channel message: ", jsException), kj::none); |
| 49 | }); |
| 50 | } |
| 51 | } |
| 52 | |
| 53 | void Channel::subscribe(jsg::Lock& js, jsg::Identified<MessageCallback> callback) { |
| 54 | subscribers.upsert(kj::mv(callback.identity), kj::mv(callback.unwrapped), [&](auto&, auto&&) {}); |
| 55 | } |
| 56 | |
| 57 | void Channel::unsubscribe(jsg::Lock& js, jsg::Identified<MessageCallback> callback) { |
| 58 | subscribers.erase(callback.identity); |
| 59 | } |
| 60 | |
| 61 | void Channel::bindStore(jsg::Lock& js, |
| 62 | jsg::Ref<AsyncLocalStorage> als, |
| 63 | jsg::Optional<TransformCallback> maybeTransform) { |
| 64 | auto key = als->getKey(); |
| 65 | KJ_IF_SOME(entry, stores.find(*key)) { |
| 66 | KJ_IF_SOME(transform, maybeTransform) { |
| 67 | entry.transform = kj::mv(transform); |
| 68 | } else { |
| 69 | entry.transform = [](jsg::Lock& js, jsg::Value value) { |
| 70 | return identityTransform(js, kj::mv(value)); |
| 71 | }; |
| 72 | } |
| 73 | return; |
| 74 | } |
| 75 | |
| 76 | KJ_IF_SOME(transform, maybeTransform) { |
| 77 | stores.insert({.key = kj::mv(key), .transform = kj::mv(transform)}); |
| 78 | } else { |
| 79 | stores.insert({.key = kj::mv(key), .transform = [](jsg::Lock& js, jsg::Value value) { |
| 80 | return identityTransform(js, kj::mv(value)); |
| 81 | }}); |
| 82 | } |
| 83 | } |
| 84 | |
| 85 | void Channel::unbindStore(jsg::Lock& js, jsg::Ref<AsyncLocalStorage> als) { |
| 86 | auto key = als->getKey(); |
| 87 | stores.eraseMatch(*key); |
| 88 | } |
| 89 | |
| 90 | v8::Local<v8::Value> Channel::runStores(jsg::Lock& js, |
| 91 | jsg::Value message, |
| 92 | jsg::Function<v8::Local<v8::Value>(jsg::Arguments<jsg::Value>)> callback, |
| 93 | jsg::Optional<v8::Local<v8::Value>> maybeReceiver, |
| 94 | jsg::Arguments<jsg::Value> args) { |
| 95 | struct StoreSnapshot { |
| 96 | kj::Own<jsg::AsyncContextFrame::StorageKey> key; |
| 97 | TransformCallback transform; |
| 98 | }; |
| 99 | auto snapshot = KJ_MAP(store, stores) -> StoreSnapshot { |
| 100 | return {kj::addRef(*store.key), store.transform.addRef(js)}; |
| 101 | }; |
| 102 | kj::Vector<kj::Own<jsg::AsyncContextFrame::StorageScope>> storageScopes; |
| 103 | for (auto& entry: snapshot) { |
| 104 | storageScopes.add(kj::heap<jsg::AsyncContextFrame::StorageScope>( |
| 105 | js, *entry.key, entry.transform(js, message.addRef(js)))); |
| 106 | } |
| 107 | |
| 108 | publish(js, message.addRef(js)); |
| 109 | |
| 110 | v8::Local<v8::Value> receiver = js.v8Context()->Global(); |
| 111 | KJ_IF_SOME(val, maybeReceiver) { |
| 112 | receiver = val; |
| 113 | } |
| 114 | callback.setReceiver(js.v8Ref(receiver)); |
| 115 | return callback(js, kj::mv(args)); |
| 116 | } |
| 117 | |
| 118 | void Channel::visitForGc(jsg::GcVisitor& visitor) { |
| 119 | for (auto& sub: subscribers) { |
| 120 | visitor.visit(sub.key, sub.value); |
| 121 | } |
| 122 | for (auto& store: stores) { |
| 123 | visitor.visit(store.transform); |
| 124 | } |
| 125 | } |
| 126 | |
| 127 | bool DiagnosticsChannelModule::hasSubscribers(jsg::Lock& js, jsg::Name name) { |
| 128 | return tryGetChannel(js, name).map([&](Channel& channel) { |
| 129 | return channel.hasSubscribers(); |
| 130 | }).orDefault(false); |
| 131 | } |
| 132 | |
| 133 | void DiagnosticsChannelModule::subscribe( |
| 134 | jsg::Lock& js, jsg::Name name, jsg::Identified<Channel::MessageCallback> callback) { |
| 135 | channel(js, kj::mv(name))->subscribe(js, kj::mv(callback)); |
| 136 | } |
| 137 | |
| 138 | void DiagnosticsChannelModule::unsubscribe( |
| 139 | jsg::Lock& js, jsg::Name name, jsg::Identified<Channel::MessageCallback> callback) { |
| 140 | KJ_IF_SOME(channel, tryGetChannel(js, name)) { |
| 141 | channel.unsubscribe(js, kj::mv(callback)); |
| 142 | } |
| 143 | } |
| 144 | |
| 145 | jsg::Ref<Channel> DiagnosticsChannelModule::channel(jsg::Lock& js, jsg::Name channel) { |
| 146 | kj::String name = channel.toString(js); |
| 147 | return channels |
| 148 | .findOrCreate(name, |
| 149 | [&, channel = kj::mv(channel)]() mutable |
| 150 | -> kj::HashMap<kj::String, jsg::Ref<Channel>>::Entry { |
| 151 | return {kj::mv(name), js.alloc<Channel>(kj::mv(channel))}; |
| 152 | }).addRef(); |
| 153 | } |
| 154 | |
| 155 | kj::Maybe<Channel&> DiagnosticsChannelModule::tryGetChannel(jsg::Lock& js, jsg::Name& name) { |
| 156 | return channels.find(name.toString(js)).map([](jsg::Ref<Channel>& channel) -> Channel& { |
| 157 | return *channel; |
| 158 | }); |
| 159 | } |
| 160 | |
| 161 | void DiagnosticsChannelModule::visitForGc(jsg::GcVisitor& visitor) { |
| 162 | for (auto& channel: channels) { |
| 163 | visitor.visit(channel.value); |
| 164 | } |
| 165 | } |
| 166 | |
| 167 | void Channel::visitForMemoryInfo(jsg::MemoryTracker& tracker) const { |
| 168 | tracker.trackField("name", name); |
| 169 | for (auto& sub: subscribers) { |
| 170 | tracker.trackField("subscribers", sub.key); |
| 171 | tracker.trackField("subscribers", sub.value); |
| 172 | } |
| 173 | for (auto& store: stores) { |
| 174 | tracker.trackField("stores", store); |
| 175 | } |
| 176 | } |
| 177 | |
| 178 | void DiagnosticsChannelModule::visitForMemoryInfo(jsg::MemoryTracker& tracker) const { |
| 179 | for (auto& channel: channels) { |
| 180 | tracker.trackField(nullptr, channel.value); |
| 181 | } |
| 182 | } |
| 183 | |
| 184 | } // namespace workerd::api::node |