Skip to content
File

Blob: src/workerd/api/node/diagnostics-channel.h

cpp117 lines
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 
10namespace workerd::api::node {
11 
12class 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 
83class 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