File
Blob: src/workerd/api/node/diagnostics-channel.h
| 1 | #pragma once |
| 2 | |
| 3 | #include <workerd/api/node/async-hooks.h> |
| 4 | #include <workerd/io/compatibility-date.capnp.h> |
| 5 | #include <workerd/jsg/jsg.h> |
| 6 | |
| 7 | #include <kj/map.h> |
| 8 | #include <kj/table.h> |
| 9 | |
| 10 | namespace workerd::api::node { |
| 11 | |
| 12 | class Channel: public jsg::Object { |
| 13 | public: |
| 14 | using MessageCallback = jsg::Function<void(jsg::Value, jsg::Name)>; |
| 15 | using TransformCallback = jsg::Function<jsg::Value(jsg::Value)>; |
| 16 | |
| 17 | static jsg::Value identityTransform(jsg::Lock& js, jsg::Value value); |
| 18 | |
| 19 | Channel(jsg::Name name); |
| 20 | |
| 21 | bool hasSubscribers(); |
| 22 | void publish(jsg::Lock& js, jsg::Value message); |
| 23 | void subscribe(jsg::Lock& js, jsg::Identified<MessageCallback> callback); |
| 24 | void unsubscribe(jsg::Lock& js, jsg::Identified<MessageCallback> callback); |
| 25 | void bindStore(jsg::Lock& js, |
| 26 | jsg::Ref<AsyncLocalStorage> als, |
| 27 | jsg::Optional<TransformCallback> maybeTransform); |
| 28 | void unbindStore(jsg::Lock& js, jsg::Ref<AsyncLocalStorage> als); |
| 29 | v8::Local<v8::Value> runStores(jsg::Lock& js, |
| 30 | jsg::Value message, |
| 31 | jsg::Function<v8::Local<v8::Value>(jsg::Arguments<jsg::Value>)> callback, |
| 32 | jsg::Optional<v8::Local<v8::Value>> maybeReceiver, |
| 33 | jsg::Arguments<jsg::Value> args); |
| 34 | |
| 35 | JSG_RESOURCE_TYPE(Channel, CompatibilityFlags::Reader flags) { |
| 36 | if (flags.getDiagnosticsChannelHasSubscribersGetter()) { |
| 37 | JSG_READONLY_PROTOTYPE_PROPERTY(hasSubscribers, hasSubscribers); |
| 38 | } else { |
| 39 | JSG_METHOD(hasSubscribers); |
| 40 | } |
| 41 | JSG_METHOD(publish); |
| 42 | JSG_METHOD(subscribe); |
| 43 | JSG_METHOD(unsubscribe); |
| 44 | JSG_METHOD(bindStore); |
| 45 | JSG_METHOD(unbindStore); |
| 46 | JSG_METHOD(runStores); |
| 47 | } |
| 48 | |
| 49 | const jsg::Name& getName() const; |
| 50 | |
| 51 | void visitForMemoryInfo(jsg::MemoryTracker& tracker) const; |
| 52 | |
| 53 | private: |
| 54 | struct StoreEntry { |
| 55 | kj::Own<jsg::AsyncContextFrame::StorageKey> key; |
| 56 | TransformCallback transform; |
| 57 | JSG_MEMORY_INFO(StoreEntry) { |
| 58 | tracker.trackField("transform", transform); |
| 59 | } |
| 60 | }; |
| 61 | |
| 62 | struct StoreCallbacks { |
| 63 | auto& keyForRow(StoreEntry& row) const { |
| 64 | return *row.key; |
| 65 | } |
| 66 | |
| 67 | bool matches(const StoreEntry& a, jsg::AsyncContextFrame::StorageKey& key) const { |
| 68 | return a.key.get() == &key; |
| 69 | } |
| 70 | |
| 71 | uint hashCode(jsg::AsyncContextFrame::StorageKey& key) const { |
| 72 | return key.hashCode(); |
| 73 | } |
| 74 | }; |
| 75 | |
| 76 | jsg::Name name; |
| 77 | kj::HashMap<jsg::HashableV8Ref<v8::Object>, MessageCallback> subscribers; |
| 78 | kj::Table<StoreEntry, kj::HashIndex<StoreCallbacks>> stores; |
| 79 | |
| 80 | void visitForGc(jsg::GcVisitor& visitor); |
| 81 | }; |
| 82 | |
| 83 | class DiagnosticsChannelModule: public jsg::Object { |
| 84 | public: |
| 85 | DiagnosticsChannelModule() = default; |
| 86 | DiagnosticsChannelModule(jsg::Lock&, const jsg::Url&) {} |
| 87 | |
| 88 | bool hasSubscribers(jsg::Lock& js, jsg::Name name); |
| 89 | jsg::Ref<Channel> channel(jsg::Lock& js, jsg::Name name); |
| 90 | void subscribe(jsg::Lock& js, jsg::Name name, jsg::Identified<Channel::MessageCallback> callback); |
| 91 | void unsubscribe( |
| 92 | jsg::Lock& js, jsg::Name name, jsg::Identified<Channel::MessageCallback> callback); |
| 93 | // TODO: Support tracing channels |
| 94 | |
| 95 | JSG_RESOURCE_TYPE(DiagnosticsChannelModule) { |
| 96 | JSG_METHOD(hasSubscribers); |
| 97 | JSG_METHOD(channel); |
| 98 | JSG_METHOD(subscribe); |
| 99 | JSG_METHOD(unsubscribe); |
| 100 | JSG_NESTED_TYPE(Channel); |
| 101 | } |
| 102 | |
| 103 | kj::Maybe<Channel&> tryGetChannel(jsg::Lock& js, jsg::Name& name); |
| 104 | |
| 105 | void visitForMemoryInfo(jsg::MemoryTracker& tracker) const; |
| 106 | |
| 107 | private: |
| 108 | kj::HashMap<kj::String, jsg::Ref<Channel>> channels; |
| 109 | |
| 110 | void visitForGc(jsg::GcVisitor& visitor); |
| 111 | }; |
| 112 | |
| 113 | #define EW_NODE_DIAGNOSTICCHANNEL_ISOLATE_TYPES \ |
| 114 | api::node::Channel, api::node::DiagnosticsChannelModule |
| 115 | |
| 116 | } // namespace workerd::api::node |