File
Blob: src/cloudflare/internal/test/instrumentation-test-helper.js
| 1 | // Copyright (c) 2025 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 | // TODO(o11y): Refactor to remove redundant code, and merge with |
| 6 | // src/workerd/api/tests/instrumentation-tail-worker.js |
| 7 | |
| 8 | import * as assert from 'node:assert'; |
| 9 | |
| 10 | /** |
| 11 | * Common helper functions for instrumentation tests. |
| 12 | * This module provides utilities for collecting and processing spans |
| 13 | * from streaming tail workers during tests. |
| 14 | */ |
| 15 | |
| 16 | /** |
| 17 | * Creates module-level state for instrumentation tests. |
| 18 | * This mirrors the original test pattern with module-level variables. |
| 19 | * @returns {Object} State object with invocationPromises and spans |
| 20 | */ |
| 21 | export function createInstrumentationState() { |
| 22 | return { |
| 23 | invocationPromises: [], |
| 24 | spans: new Map(), |
| 25 | }; |
| 26 | } |
| 27 | |
| 28 | /** |
| 29 | * Creates the tailStream handler for instrumentation tests. |
| 30 | * @param {Object} state - The state object from createInstrumentationState |
| 31 | * @returns {Function} The tailStream handler function |
| 32 | */ |
| 33 | export function createTailStreamHandler(state) { |
| 34 | return (event, env, ctx) => { |
| 35 | // For each "onset" event, store a promise which we will resolve when |
| 36 | // we receive the equivalent "outcome" event |
| 37 | let resolveFn; |
| 38 | state.invocationPromises.push( |
| 39 | new Promise((resolve, reject) => { |
| 40 | resolveFn = resolve; |
| 41 | }) |
| 42 | ); |
| 43 | |
| 44 | // Accumulate the span info for easier testing |
| 45 | return (event) => { |
| 46 | let spanKey = `${event.invocationId}#${event.event.spanId || event.spanContext.spanId}`; |
| 47 | switch (event.event.type) { |
| 48 | case 'spanOpen': |
| 49 | // The span ids will change between tests, but Map preserves insertion order |
| 50 | state.spans.set(spanKey, { name: event.event.name }); |
| 51 | break; |
| 52 | case 'attributes': { |
| 53 | let span = state.spans.get(spanKey); |
| 54 | for (let { name, value } of event.event.info) { |
| 55 | span[name] = value; |
| 56 | } |
| 57 | state.spans.set(spanKey, span); |
| 58 | break; |
| 59 | } |
| 60 | case 'spanClose': { |
| 61 | let span = state.spans.get(spanKey); |
| 62 | span['closed'] = true; |
| 63 | state.spans.set(spanKey, span); |
| 64 | break; |
| 65 | } |
| 66 | case 'outcome': |
| 67 | resolveFn(); |
| 68 | break; |
| 69 | } |
| 70 | }; |
| 71 | }; |
| 72 | } |
| 73 | |
| 74 | /** |
| 75 | * Creates a tail-stream handler that records span parent/child relationships in addition |
| 76 | * to the span attributes. Use this when your test needs to assert on hierarchy (e.g., |
| 77 | * "the fetch span's parent is the enterSpan I just opened"). |
| 78 | * |
| 79 | * The returned state shape is: |
| 80 | * state.spans: Map<spanKey, { |
| 81 | * name, spanId, parentSpanId, invocationId, closed, ...attributes |
| 82 | * }> |
| 83 | * state.topLevelSpans: Map<invocationId, topLevelSpanId> |
| 84 | * (the span ID from the onset event, i.e. the implicit request root) |
| 85 | * |
| 86 | * @returns {{state: Object, tailStream: Function, waitForCompletion: Function}} |
| 87 | */ |
| 88 | export function createHierarchyAwareCollector() { |
| 89 | const state = { |
| 90 | invocationPromises: [], |
| 91 | spans: new Map(), |
| 92 | topLevelSpans: new Map(), |
| 93 | }; |
| 94 | |
| 95 | const tailStream = (event, env, ctx) => { |
| 96 | // Record the onset event's spanId as the top-level span for this invocation. |
| 97 | state.topLevelSpans.set(event.invocationId, event.event.spanId); |
| 98 | |
| 99 | let resolveFn; |
| 100 | state.invocationPromises.push( |
| 101 | new Promise((resolve) => { |
| 102 | resolveFn = resolve; |
| 103 | }) |
| 104 | ); |
| 105 | |
| 106 | return (event) => { |
| 107 | const spanKey = `${event.invocationId}#${event.event.spanId || event.spanContext.spanId}`; |
| 108 | switch (event.event.type) { |
| 109 | case 'spanOpen': |
| 110 | // event.event.spanId is the new span's id; event.spanContext.spanId is the |
| 111 | // span that was active when this one was opened (i.e., its parent). |
| 112 | state.spans.set(spanKey, { |
| 113 | name: event.event.name, |
| 114 | spanId: event.event.spanId, |
| 115 | parentSpanId: event.spanContext.spanId, |
| 116 | invocationId: event.invocationId, |
| 117 | }); |
| 118 | break; |
| 119 | case 'attributes': { |
| 120 | const span = state.spans.get(spanKey); |
| 121 | if (!span) break; |
| 122 | for (const { name, value } of event.event.info) { |
| 123 | span[name] = value; |
| 124 | } |
| 125 | break; |
| 126 | } |
| 127 | case 'spanClose': { |
| 128 | const span = state.spans.get(spanKey); |
| 129 | if (!span) break; |
| 130 | span.closed = true; |
| 131 | break; |
| 132 | } |
| 133 | case 'outcome': |
| 134 | resolveFn(); |
| 135 | break; |
| 136 | } |
| 137 | }; |
| 138 | }; |
| 139 | |
| 140 | const waitForCompletion = () => Promise.allSettled(state.invocationPromises); |
| 141 | |
| 142 | return { state, tailStream, waitForCompletion }; |
| 143 | } |
| 144 | |
| 145 | /** |
| 146 | * Find a span whose `name` field matches (after an optional filter function). |
| 147 | * Throws if there is not exactly one such span in the collector's state. |
| 148 | */ |
| 149 | export function findSpanByName(state, name, filterFn = () => true) { |
| 150 | const matches = []; |
| 151 | for (const span of state.spans.values()) { |
| 152 | if (span.name === name && filterFn(span)) matches.push(span); |
| 153 | } |
| 154 | if (matches.length === 0) { |
| 155 | throw new Error(`No span found with name='${name}'`); |
| 156 | } |
| 157 | if (matches.length > 1) { |
| 158 | throw new Error( |
| 159 | `Expected exactly one span with name='${name}', found ${matches.length}: ` + |
| 160 | JSON.stringify(matches, null, 2) |
| 161 | ); |
| 162 | } |
| 163 | return matches[0]; |
| 164 | } |
| 165 | |
| 166 | /** |
| 167 | * Runs instrumentation test assertions. |
| 168 | * This mirrors the original test logic exactly. |
| 169 | * @param {Object} state - The state object from createInstrumentationState |
| 170 | * @param {Array} expectedSpans - The expected spans to compare against |
| 171 | * @param {Object} options - Options for the test |
| 172 | * @param {Function} options.mapFn - Map function to transform spans before comparison (default: x => x) |
| 173 | * @param {Function} options.filterFn - Filter function for spans (default: filters out jsRpcSession) |
| 174 | * @param {string} options.testName - Name for the test (default: 'instrumentation') |
| 175 | * @param {boolean} options.logReceived - Log received spans for debugging (default: false) |
| 176 | * |
| 177 | * Usage for updating tests: |
| 178 | * await runInstrumentationTest(state, expectedSpans, { logReceived: true }); |
| 179 | * // Copy the logged output to update expectedSpans |
| 180 | */ |
| 181 | export async function runInstrumentationTest( |
| 182 | state, |
| 183 | expectedSpans, |
| 184 | options = {} |
| 185 | ) { |
| 186 | const { |
| 187 | mapFn = (x) => x, |
| 188 | filterFn = (span) => span.name !== 'jsRpcSession', |
| 189 | testName = 'instrumentation', |
| 190 | logReceived = false, |
| 191 | } = options; |
| 192 | |
| 193 | // Wait for all the tailStream executions to finish |
| 194 | await Promise.allSettled(state.invocationPromises); |
| 195 | |
| 196 | // Recorded streaming tail worker events, in insertion order, |
| 197 | // mapping and filtering spans not associated with the test |
| 198 | let received = Array.from(state.spans.values()).map(mapFn).filter(filterFn); |
| 199 | |
| 200 | // Log received spans for debugging/updating tests |
| 201 | if (logReceived) { |
| 202 | console.log(`Received spans for ${testName}:\n`, received); |
| 203 | } |
| 204 | |
| 205 | let failed = 0; |
| 206 | let i = -1; |
| 207 | |
| 208 | try { |
| 209 | assert.equal(received.length, expectedSpans.length); |
| 210 | for (i = 0; i < received.length; i++) { |
| 211 | assert.deepStrictEqual(received[i], expectedSpans[i]); |
| 212 | } |
| 213 | } catch (e) { |
| 214 | failed++; |
| 215 | if (i >= 0) { |
| 216 | console.log('spans are not identical', e); |
| 217 | } else { |
| 218 | console.error(e); |
| 219 | } |
| 220 | } |
| 221 | |
| 222 | if (failed > 0) { |
| 223 | throw `${testName} test failed`; |
| 224 | } |
| 225 | } |
| 226 | |
| 227 | /** |
| 228 | * Creates a tail stream collector for instrumentation tests with encapsulated state. |
| 229 | * This provides a different API style where state is encapsulated in the returned object. |
| 230 | * @returns {Object} An object with methods to handle spans |
| 231 | */ |
| 232 | export function createTailStreamCollector() { |
| 233 | let state = createInstrumentationState(); |
| 234 | |
| 235 | const tailStream = createTailStreamHandler(state); |
| 236 | |
| 237 | let spans = state.spans; |
| 238 | let invocationPromises = state.invocationPromises; |
| 239 | const waitForCompletion = () => { |
| 240 | return Promise.allSettled(invocationPromises); |
| 241 | }; |
| 242 | |
| 243 | return { |
| 244 | tailStream, |
| 245 | waitForCompletion, |
| 246 | spans, |
| 247 | }; |
| 248 | } |
| 249 | |
| 250 | /** |
| 251 | * Groups spans by a specific attribute. |
| 252 | * @param {Array} spans - The spans to group |
| 253 | * @param {string} attribute - The attribute to group by |
| 254 | * @returns {Map} A map of attribute values to arrays of spans |
| 255 | */ |
| 256 | export function groupSpansBy(spans, attribute) { |
| 257 | const groups = new Map(); |
| 258 | |
| 259 | for (const span of spans) { |
| 260 | const key = span[attribute] || 'unknown'; |
| 261 | if (!groups.has(key)) { |
| 262 | groups.set(key, []); |
| 263 | } |
| 264 | groups.get(key).push(span); |
| 265 | } |
| 266 | |
| 267 | return groups; |
| 268 | } |