Skip to content
File

Blob: src/workerd/api/node/tests/diagnostics-channel-test.js

javascript232 lines
1// Copyright (c) 2023 Cloudflare, Inc.
2// Licensed under the Apache 2.0 license found in the LICENSE file or at:
3// https://opensource.org/licenses/Apache-2.0
4import { ok, strictEqual } from 'node:assert';
5 
6import {
7 hasSubscribers,
8 channel,
9 subscribe,
10 unsubscribe,
11 tracingChannel,
12 Channel,
13} from 'node:diagnostics_channel';
14 
15import { AsyncLocalStorage } from 'node:async_hooks';
16 
17export const test_basics = {
18 async test(ctrl, env, ctx) {
19 ok(!hasSubscribers('foo'));
20 const channel1 = channel('foo');
21 const channel2 = channel('foo');
22 strictEqual(channel1, channel2);
23 ok(channel1 instanceof Channel);
24 
25 const messagePromise = Promise.withResolvers();
26 
27 const listener = (message) => {
28 try {
29 strictEqual(message, 'hello');
30 messagePromise.resolve();
31 } catch (err) {
32 messagePromise.reject(err);
33 }
34 };
35 
36 subscribe('foo', listener);
37 
38 ok(hasSubscribers('foo'));
39 
40 channel1.publish('hello');
41 
42 unsubscribe('foo', listener);
43 
44 ok(!hasSubscribers('foo'));
45 
46 await messagePromise.promise;
47 },
48};
49 
50export const test_tracing = {
51 async test(ctrl, env, ctx) {
52 const tc = tracingChannel('bar');
53 ok(tc.start instanceof Channel);
54 ok(tc.end instanceof Channel);
55 ok(tc.asyncStart instanceof Channel);
56 ok(tc.asyncEnd instanceof Channel);
57 ok(tc.error instanceof Channel);
58 
59 const als = new AsyncLocalStorage();
60 tc.start.bindStore(als);
61 
62 const promises = [
63 Promise.withResolvers(),
64 Promise.withResolvers(),
65 Promise.withResolvers(),
66 Promise.withResolvers(),
67 Promise.withResolvers(),
68 ];
69 
70 const context = {};
71 
72 tc.subscribe({
73 start(_, name) {
74 try {
75 // Since the als is bound to tc.start, the context should be
76 // propagated here to the listener.
77 strictEqual(als.getStore(), context);
78 strictEqual(name, 'tracing:bar:start');
79 promises[0].resolve();
80 } catch (err) {
81 promises[0].reject(err);
82 }
83 },
84 end(_, name) {
85 try {
86 // Since the als is bound to tc.start, the context should be
87 // propagated here to the other listeners even if they aren't
88 // explicitly bound to als.
89 strictEqual(als.getStore(), context);
90 strictEqual(name, 'tracing:bar:end');
91 promises[1].resolve();
92 } catch (err) {
93 promises[1].reject(err);
94 }
95 },
96 asyncStart() {
97 promises[2].resolve();
98 },
99 asyncEnd() {
100 promises[3].resolve();
101 },
102 error() {
103 promises[4].resolve();
104 },
105 });
106 
107 tc.tracePromise(async () => {
108 throw new Error('boom');
109 }, context);
110 
111 await Promise.all(promises.map((p) => p.promise));
112 },
113};
114 
115export const serFailureTest = {
116 async test() {
117 const { promise, resolve } = Promise.withResolvers();
118 const channel1 = channel('ser');
119 subscribe('ser', () => {
120 resolve();
121 });
122 channel1.publish(function () {});
123 await promise;
124 },
125};
126 
127export const DiagChannelSubscribeDuringPublish = {
128 async test() {
129 const ch = channel('test-publish');
130 let callCount = 0;
131 let duringPublishCallCount = 0;
132 
133 for (let i = 0; i < 32; i++) {
134 ch.subscribe(() => {
135 callCount++;
136 });
137 }
138 
139 ch.subscribe(() => {
140 callCount++;
141 for (let i = 0; i < 64; i++) {
142 ch.subscribe(() => {
143 duringPublishCallCount++;
144 });
145 }
146 });
147 
148 for (let i = 0; i < 16; i++) {
149 ch.subscribe(() => {
150 callCount++;
151 });
152 }
153 
154 ch.publish({ data: 'trigger' });
155 
156 strictEqual(callCount, 49);
157 strictEqual(duringPublishCallCount, 0);
158 },
159};
160 
161export const DiagChannelUnbindDuringRunStores = {
162 async test() {
163 const ch = channel('test');
164 const als = new AsyncLocalStorage();
165 let transformCallCount = 0;
166 
167 ch.bindStore(als, (msg) => {
168 transformCallCount++;
169 ch.unbindStore(als);
170 return msg;
171 });
172 
173 const result = ch.runStores({}, () => 'done');
174 strictEqual(result, 'done');
175 strictEqual(transformCallCount, 1);
176 
177 ch.runStores({}, () => {});
178 strictEqual(transformCallCount, 1);
179 },
180};
181 
182// Legacy-method behavior: when the `diagnostics_channel_has_subscribers_getter`
183// compat flag is NOT enabled, `Channel.hasSubscribers` and
184// `TracingChannel.hasSubscribers` are registered as methods. The associated
185// `.wd-test` opts out of the flag explicitly so that this behavior is exercised
186// under both the default and `@all-compat-flags` test variants.
187export const test_channel_hasSubscribers_is_a_method = {
188 async test() {
189 const ch = channel('method-test');
190 
191 // It is a method, not a boolean property.
192 strictEqual(typeof ch.hasSubscribers, 'function');
193 strictEqual(ch.hasSubscribers(), false);
194 
195 const listener = () => {};
196 ch.subscribe(listener);
197 strictEqual(ch.hasSubscribers(), true);
198 
199 ch.unsubscribe(listener);
200 strictEqual(ch.hasSubscribers(), false);
201 
202 // Defined on the prototype as a method (value descriptor, no getter).
203 const desc = Object.getOwnPropertyDescriptor(
204 Object.getPrototypeOf(ch),
205 'hasSubscribers'
206 );
207 ok(desc);
208 strictEqual(typeof desc.value, 'function');
209 strictEqual(desc.get, undefined);
210 strictEqual(desc.set, undefined);
211 },
212};
213 
214export const test_tracingChannel_hasSubscribers_is_a_method = {
215 async test() {
216 const tc = tracingChannel('tracing-method-test');
217 
218 strictEqual(typeof tc.hasSubscribers, 'function');
219 strictEqual(tc.hasSubscribers(), false);
220 
221 const listener = () => {};
222 
223 // Each sub-channel independently flips the aggregate method result.
224 for (const sub of ['start', 'end', 'asyncStart', 'asyncEnd', 'error']) {
225 tc[sub].subscribe(listener);
226 strictEqual(tc.hasSubscribers(), true, `via ${sub}`);
227 tc[sub].unsubscribe(listener);
228 strictEqual(tc.hasSubscribers(), false, `after unsubscribing ${sub}`);
229 }
230 },
231};