Skip to content
File

Blob: src/node/diagnostics_channel.ts

typescript358 lines
1// Copyright (c) 2017-2022 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
4//
5// Copyright Joyent, Inc. and other Node contributors.
6//
7// Permission is hereby granted, free of charge, to any person obtaining a
8// copy of this software and associated documentation files (the
9// "Software"), to deal in the Software without restriction, including
10// without limitation the rights to use, copy, modify, merge, publish,
11// distribute, sublicense, and/or sell copies of the Software, and to permit
12// persons to whom the Software is furnished to do so, subject to the
13// following conditions:
14//
15// The above copyright notice and this permission notice shall be included
16// in all copies or substantial portions of the Software.
17//
18// THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS
19// OR IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF
20// MERCHANTABILITY, FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN
21// NO EVENT SHALL THE AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM,
22// DAMAGES OR OTHER LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR
23// OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE
24// USE OR OTHER DEALINGS IN THE SOFTWARE.
25 
26import { default as diagnosticsChannel } from 'node-internal:diagnostics_channel';
27 
28import type {
29 Channel as ChannelType,
30 MessageCallback,
31} from 'node-internal:diagnostics_channel';
32 
33import { ERR_INVALID_ARG_TYPE } from 'node-internal:internal_errors';
34 
35import { validateObject } from 'node-internal:validators';
36 
37const hasSubscribersGetter =
38 !!Cloudflare.compatibilityFlags['diagnostics_channel_has_subscribers_getter'];
39 
40export const { Channel } = diagnosticsChannel;
41 
42export function hasSubscribers(name: string | symbol): boolean {
43 return diagnosticsChannel.hasSubscribers(name);
44}
45 
46export function channel(name: string | symbol): ChannelType {
47 return diagnosticsChannel.channel(name);
48}
49 
50export function subscribe(
51 name: string | symbol,
52 callback: MessageCallback
53): void {
54 diagnosticsChannel.subscribe(name, callback);
55}
56 
57export function unsubscribe(
58 name: string | symbol,
59 callback: MessageCallback
60): void {
61 diagnosticsChannel.unsubscribe(name, callback);
62}
63 
64export interface TracingChannelSubscriptions {
65 start?: MessageCallback;
66 end?: MessageCallback;
67 asyncStart?: MessageCallback;
68 asyncEnd?: MessageCallback;
69 error?: MessageCallback;
70}
71 
72export interface TracingChannels {
73 start: ChannelType;
74 end: ChannelType;
75 asyncStart: ChannelType;
76 asyncEnd: ChannelType;
77 error: ChannelType;
78}
79 
80const kStart = Symbol('kStart');
81const kEnd = Symbol('kEnd');
82const kAsyncStart = Symbol('kAsyncStart');
83const kAsyncEnd = Symbol('kAsyncEnd');
84const kError = Symbol('kError');
85 
86export class TracingChannel {
87 private [kStart]?: ChannelType;
88 private [kEnd]?: ChannelType;
89 private [kAsyncStart]?: ChannelType;
90 private [kAsyncEnd]?: ChannelType;
91 private [kError]?: ChannelType;
92 
93 constructor() {
94 throw new Error(
95 'Use diagnostic_channel.tracingChannels() to create TracingChannel'
96 );
97 }
98 
99 get start(): ChannelType | undefined {
100 return this[kStart];
101 }
102 get end(): ChannelType | undefined {
103 return this[kEnd];
104 }
105 get asyncStart(): ChannelType | undefined {
106 return this[kAsyncStart];
107 }
108 get asyncEnd(): ChannelType | undefined {
109 return this[kAsyncEnd];
110 }
111 get error(): ChannelType | undefined {
112 return this[kError];
113 }
114 
115 // `hasSubscribers` is installed on the prototype below, gated by the
116 // `diagnostics_channel_has_subscribers_getter` compatibility flag. When the
117 // flag is enabled it's a getter property (matching Node.js); when disabled
118 // it's a method (preserving the original workerd behavior).
119 static {
120 if (hasSubscribersGetter) {
121 Object.defineProperty(this.prototype, 'hasSubscribers', {
122 get(this: TracingChannel): boolean {
123 return [
124 this.start,
125 this.end,
126 this.asyncStart,
127 this.asyncEnd,
128 this.error,
129 ].some((c) => c != null && (c.hasSubscribers as unknown as boolean));
130 },
131 configurable: true,
132 enumerable: false,
133 });
134 } else {
135 Object.defineProperty(this.prototype, 'hasSubscribers', {
136 value: function hasSubscribers(this: TracingChannel): boolean {
137 return [
138 this.start,
139 this.end,
140 this.asyncStart,
141 this.asyncEnd,
142 this.error,
143 ].some(
144 (c) =>
145 c != null &&
146 (c.hasSubscribers as unknown as () => boolean).call(c)
147 );
148 },
149 writable: true,
150 configurable: true,
151 enumerable: false,
152 });
153 }
154 }
155 
156 subscribe(subscriptions: TracingChannelSubscriptions): void {
157 if (subscriptions.start !== undefined)
158 this[kStart]?.subscribe(subscriptions.start);
159 if (subscriptions.end !== undefined)
160 this[kEnd]?.subscribe(subscriptions.end);
161 if (subscriptions.asyncStart !== undefined)
162 this[kAsyncStart]?.subscribe(subscriptions.asyncStart);
163 if (subscriptions.asyncEnd !== undefined)
164 this[kAsyncEnd]?.subscribe(subscriptions.asyncEnd);
165 if (subscriptions.error !== undefined)
166 this[kError]?.subscribe(subscriptions.error);
167 }
168 
169 unsubscribe(subscriptions: TracingChannelSubscriptions): void {
170 if (subscriptions.start !== undefined)
171 this[kStart]?.unsubscribe(subscriptions.start);
172 if (subscriptions.end !== undefined)
173 this[kEnd]?.unsubscribe(subscriptions.end);
174 if (subscriptions.asyncStart !== undefined)
175 this[kAsyncStart]?.unsubscribe(subscriptions.asyncStart);
176 if (subscriptions.asyncEnd !== undefined)
177 this[kAsyncEnd]?.unsubscribe(subscriptions.asyncEnd);
178 if (subscriptions.error !== undefined)
179 this[kError]?.unsubscribe(subscriptions.error);
180 }
181 
182 traceSync(
183 fn: (...args: unknown[]) => unknown,
184 context: Record<string, unknown> = {},
185 thisArg: unknown = globalThis,
186 ...args: unknown[]
187 ): unknown {
188 const { start, end, error } = this;
189 
190 return start?.runStores(
191 context,
192 () => {
193 try {
194 const result = Reflect.apply(fn, thisArg, args);
195 context.result = result;
196 return result;
197 } catch (err) {
198 context.error = err;
199 error?.publish(context);
200 throw err;
201 } finally {
202 end?.publish(context);
203 }
204 },
205 thisArg
206 );
207 }
208 
209 tracePromise(
210 fn: (...args: unknown[]) => unknown,
211 context: Record<string, unknown> = {},
212 thisArg: unknown = globalThis,
213 ...args: unknown[]
214 ): unknown {
215 const { start, end, asyncStart, asyncEnd, error } = this;
216 
217 function reject(err: Error): Promise<unknown> {
218 context.error = err;
219 error?.publish(context);
220 asyncStart?.publish(context);
221 asyncEnd?.publish(context);
222 return Promise.reject(err);
223 }
224 
225 function resolve(result: unknown): unknown {
226 context.result = result;
227 asyncStart?.publish(context);
228 asyncEnd?.publish(context);
229 return result;
230 }
231 
232 return start?.runStores(
233 context,
234 () => {
235 try {
236 let promise = Reflect.apply(fn, thisArg, args) as Promise<unknown>;
237 // Convert thenables to native promises
238 if (!(promise instanceof Promise)) {
239 promise = Promise.resolve(promise);
240 }
241 return promise.then(resolve, reject);
242 } catch (err) {
243 context.error = err;
244 error?.publish(context);
245 throw err;
246 } finally {
247 end?.publish(context);
248 }
249 },
250 thisArg
251 );
252 }
253 
254 traceCallback(
255 fn: (...args: unknown[]) => unknown,
256 position = -1,
257 context: Record<string, unknown> = {},
258 thisArg: unknown = globalThis,
259 ...args: unknown[]
260 ): unknown {
261 const { start, end, asyncStart, asyncEnd, error } = this;
262 
263 function wrappedCallback(this: unknown, err: unknown, res: unknown): void {
264 if (err) {
265 context.error = err;
266 error?.publish(context);
267 } else {
268 context.result = res;
269 }
270 
271 // Using runStores here enables manual context failure recovery
272 asyncStart?.runStores(
273 context,
274 () => {
275 try {
276 if (callback) {
277 // eslint-disable-next-line prefer-rest-params
278 Reflect.apply(callback, this, arguments);
279 }
280 } finally {
281 asyncEnd?.publish(context);
282 }
283 },
284 thisArg
285 );
286 }
287 
288 const callback = args[position] as VoidFunction | undefined;
289 if (typeof callback !== 'function') {
290 throw new ERR_INVALID_ARG_TYPE('callback', ['function'], callback);
291 }
292 args.splice(position, 1, wrappedCallback);
293 
294 return start?.runStores(
295 context,
296 () => {
297 try {
298 return Reflect.apply(fn, thisArg, args);
299 } catch (err) {
300 context.error = err;
301 error?.publish(context);
302 throw err;
303 } finally {
304 end?.publish(context);
305 }
306 },
307 thisArg
308 );
309 }
310}
311 
312function validateChannel(channel: unknown, name: string): ChannelType {
313 if (!(channel instanceof Channel)) {
314 throw new ERR_INVALID_ARG_TYPE(name, 'Channel', channel);
315 }
316 return channel;
317}
318 
319export function tracingChannel(
320 nameOrChannels: string | TracingChannels
321): TracingChannel {
322 return Reflect.construct(
323 function (this: TracingChannel) {
324 if (typeof nameOrChannels === 'string') {
325 this[kStart] = channel(`tracing:${nameOrChannels}:start`);
326 this[kEnd] = channel(`tracing:${nameOrChannels}:end`);
327 this[kAsyncStart] = channel(`tracing:${nameOrChannels}:asyncStart`);
328 this[kAsyncEnd] = channel(`tracing:${nameOrChannels}:asyncEnd`);
329 this[kError] = channel(`tracing:${nameOrChannels}:error`);
330 } else {
331 validateObject(nameOrChannels, 'channels');
332 this[kStart] = validateChannel(nameOrChannels.start, 'channels.start');
333 this[kEnd] = validateChannel(nameOrChannels.end, 'channels.end');
334 this[kAsyncStart] = validateChannel(
335 nameOrChannels.asyncStart,
336 'channels.asyncStart'
337 );
338 this[kAsyncEnd] = validateChannel(
339 nameOrChannels.asyncEnd,
340 'channels.asyncEnd'
341 );
342 this[kError] = validateChannel(nameOrChannels.error, 'channels.error');
343 }
344 },
345 [],
346 TracingChannel
347 ) as TracingChannel;
348}
349 
350export default {
351 hasSubscribers,
352 channel,
353 subscribe,
354 unsubscribe,
355 tracingChannel,
356 Channel,
357};