#pragma once #include #include #include #include #include namespace workerd::api::node { class Channel: public jsg::Object { public: using MessageCallback = jsg::Function; using TransformCallback = jsg::Function; static jsg::Value identityTransform(jsg::Lock& js, jsg::Value value); Channel(jsg::Name name); bool hasSubscribers(); void publish(jsg::Lock& js, jsg::Value message); void subscribe(jsg::Lock& js, jsg::Identified callback); void unsubscribe(jsg::Lock& js, jsg::Identified callback); void bindStore(jsg::Lock& js, jsg::Ref als, jsg::Optional maybeTransform); void unbindStore(jsg::Lock& js, jsg::Ref als); v8::Local runStores(jsg::Lock& js, jsg::Value message, jsg::Function(jsg::Arguments)> callback, jsg::Optional> maybeReceiver, jsg::Arguments args); JSG_RESOURCE_TYPE(Channel, CompatibilityFlags::Reader flags) { if (flags.getDiagnosticsChannelHasSubscribersGetter()) { JSG_READONLY_PROTOTYPE_PROPERTY(hasSubscribers, hasSubscribers); } else { JSG_METHOD(hasSubscribers); } JSG_METHOD(publish); JSG_METHOD(subscribe); JSG_METHOD(unsubscribe); JSG_METHOD(bindStore); JSG_METHOD(unbindStore); JSG_METHOD(runStores); } const jsg::Name& getName() const; void visitForMemoryInfo(jsg::MemoryTracker& tracker) const; private: struct StoreEntry { kj::Own key; TransformCallback transform; JSG_MEMORY_INFO(StoreEntry) { tracker.trackField("transform", transform); } }; struct StoreCallbacks { auto& keyForRow(StoreEntry& row) const { return *row.key; } bool matches(const StoreEntry& a, jsg::AsyncContextFrame::StorageKey& key) const { return a.key.get() == &key; } uint hashCode(jsg::AsyncContextFrame::StorageKey& key) const { return key.hashCode(); } }; jsg::Name name; kj::HashMap, MessageCallback> subscribers; kj::Table> stores; void visitForGc(jsg::GcVisitor& visitor); }; class DiagnosticsChannelModule: public jsg::Object { public: DiagnosticsChannelModule() = default; DiagnosticsChannelModule(jsg::Lock&, const jsg::Url&) {} bool hasSubscribers(jsg::Lock& js, jsg::Name name); jsg::Ref channel(jsg::Lock& js, jsg::Name name); void subscribe(jsg::Lock& js, jsg::Name name, jsg::Identified callback); void unsubscribe( jsg::Lock& js, jsg::Name name, jsg::Identified callback); // TODO: Support tracing channels JSG_RESOURCE_TYPE(DiagnosticsChannelModule) { JSG_METHOD(hasSubscribers); JSG_METHOD(channel); JSG_METHOD(subscribe); JSG_METHOD(unsubscribe); JSG_NESTED_TYPE(Channel); } kj::Maybe tryGetChannel(jsg::Lock& js, jsg::Name& name); void visitForMemoryInfo(jsg::MemoryTracker& tracker) const; private: kj::HashMap> channels; void visitForGc(jsg::GcVisitor& visitor); }; #define EW_NODE_DIAGNOSTICCHANNEL_ISOLATE_TYPES \ api::node::Channel, api::node::DiagnosticsChannelModule } // namespace workerd::api::node